SiteWise Edge GatewayでOPC UAデータをKinesisにストリーミングする
SiteWise Edge Gatewayとカスタムのgreengrassコンポーネントを使って、OPC UAテレメトリをKinesis Data Streamsへ橋渡しします。
SiteWise Edge Gatewayは、OPC UAのデータソースをKinesis Data Streamsに橋渡しでき、工場現場のテレメトリを下流のAWSサービスで利用できるようにします。

バックエンドの構築
OPC UAサーバー
ダミーのOPC UAサーバーを作成するには、Pythonライブラリopcua-asyncioを使用します。
必要なパッケージをインストールします。
pip install asyncuaサンプルスクリプト(server-minimal.py)をEC2インスタンスにダウンロードします。
curl -OL https://raw.githubusercontent.com/FreeOpcUa/opcua-asyncio/master/examples/server-minimal.pyスクリプトを実行します。
python server-minimal.pyGreengrass V2コアデバイス
公式ドキュメントに従って、Greengrass V2コアデバイスをセットアップします。この記事では、セットアップ手順自体は扱いません。
インストール後、以下を実施します。
aws.greengrass.StreamManagerコンポーネントをデプロイします。- Kinesis Data Streamsへのデータ送信を許可するため、以下のIAMポリシーをトークン交換ロールに追加します。
{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "kinesis:PutRecord", "kinesis:PutRecords" ], "Resource": "arn:aws:kinesis:*:$YOUR_AWS_ACCOUNT_ID:stream/$KINESIS_DATA_STREAM_NAME" } ]}SiteWise Edge Gateway
SiteWise Edge Gatewayをセットアップします。
Advanced setupを選択します。

Data processing packは不要です。

Configure publisherステップはスキップします。

OPC UAサーバーのローカルエンドポイントを指定し、Greengrassストリーム名(例: SiteWise_Stream_Kinesis)を設定します。



設定内容を確認し、Createボタンを押します。

Kinesis Data Stream
OPC UAデータの送信先として、Kinesis Data Streamを作成します。

Kinesis Data Firehose
ストリーミングされたデータを永続化するため、以下の内容でKinesis Data Firehoseを設定します。
- ソース: 上記で作成したKinesis Data Stream。
- 送信先: S3バケット。

Greengrassコンポーネントの作成
GreengrassストリームからKinesis Data Streamsへデータをストリーミングするには、カスタムのGreengrassコンポーネントを作成する必要があります。Greengrassコンポーネントの開発の詳細については、公式ドキュメントを参照してください。
ディレクトリ構成
コンポーネントのファイルは、以下のような構成にします。
- kinesis_data_stream.py- recipe.yaml- requirements.txt- stream_manager_sdk.zipコンポーネントの作成
stream_manager_sdk.zipファイルには、Greengrassコンポーネントに必要なSDKを含める必要があります。詳細やサンプルコードについては、AWS Greengrass Stream Manager SDK for Pythonを参照してください。
requirements.txtを作成します。Greengrass Stream Managerとやり取りするには、stream-managerライブラリが不可欠です。
cbor2~=5.4.2stream-manager==1.1.1以下の内容でkinesis_data_stream.pyを作成します。このスクリプトは、Greengrass Stream Manager SDKを利用してKinesis Data Streamへデータをストリーミングします。
"""Script to use Greengrass Stream Manager to stream data to a Kinesis Data StreamSee also https://github.com/aws-greengrass/aws-greengrass-stream-manager-sdk-python/blob/main/samples/stream_manager_kinesis.py"""
import argparseimport asyncioimport loggingimport time
from stream_manager import ( ExportDefinition, KinesisConfig, MessageStreamDefinition, ReadMessagesOptions, ResourceNotFoundException, StrategyOnFull, StreamManagerClient,)
logging.basicConfig(level=logging.INFO)logger = logging.getLogger()
def main( stream_name: str, kinesis_stream_name: str, batch_size: int = None): try: # Create a client for the StreamManager client = StreamManagerClient()
# Try deleting the stream (if it exists) so that we have a fresh start try: client.delete_message_stream(stream_name=stream_name) except ResourceNotFoundException: pass
exports = ExportDefinition( kinesis=[KinesisConfig( identifier="KinesisExport" + stream_name, kinesis_stream_name=kinesis_stream_name, batch_size=batch_size, )] ) client.create_message_stream( MessageStreamDefinition( name=stream_name, strategy_on_full=StrategyOnFull.OverwriteOldestData, export_definition=exports ) )
while True: time.sleep(1)
except asyncio.TimeoutError: logger.exception("Timed out while executing") except Exception: logger.exception("Exception while running") finally: # Always close the client to avoid resource leaks if client: client.close()
def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser() parser.add_argument('--greengrass-stream', required=True, default='SiteWise_Stream_Kinesis') parser.add_argument('--kinesis-stream', required=True) parser.add_argument('--batch-size', required=False, type=int, default=500) return parser.parse_args()
if __name__ == '__main__': args = parse_args() logger.info(f'args: {args.__dict__}') main(args.greengrass_stream, args.kinesis_stream, args.batch_size)以下の内容でrecipe.yamlファイルを作成します。この記事では、コンポーネント名はjp.co.xyz.StreamManagerKinesisです。
このレシピには、ComponentConfiguration(14〜16行目)に以下の設定可能なパラメータが含まれています。
- GreengrassStream: Greengrassストリームの名前
- KinesisStream: OPC UAデータを受け取るKinesis Data Streamの名前
- BatchSize: データ転送のバッチサイズ(最小: 1、最大: 500)
# Replace $ArtifactsS3Bucket with your value to complete component registration.
RecipeFormatVersion: 2020-01-25
ComponentName: jp.co.xyz.StreamManagerKinesisComponentVersion: 1.0.0ComponentDescription: Streams data in Greengrass stream to a Kinesis Data Stream.ComponentPublisher: selfComponentDependencies: aws.greengrass.StreamManager: VersionRequirement: '^2.0.0'ComponentConfiguration: DefaultConfiguration: GreengrassStream: SiteWise_Stream_Kinesis KinesisStream: '' BatchSize: 100 # minimum 1, maximum 500
Manifests: - Platform: os: linux Lifecycle: Install: pip3 install --user -r {artifacts:decompressedPath}/component/requirements.txt Run: | export PYTHONPATH=$PYTHONPATH:{artifacts:decompressedPath}/stream_manager_sdk python3 {artifacts:decompressedPath}/component/kinesis_data_stream.py \ --greengrass-stream {configuration:/GreengrassStream} \ --kinesis-stream {configuration:/KinesisStream} \ --batch-size {configuration:/BatchSize} Artifacts: - URI: s3://$ArtifactsS3Bucket/artifacts/jp.co.xyz.StreamManagerKinesis/1.0.0/component.zip Unarchive: ZIPコンポーネントアーティファクトのアップロード
以下のコマンドで、コンポーネントのファイルをアーカイブし、S3バケットにアップロードします。
S3_BUCKET=<YOUR_BUCKET_NAME>VERSION=1.0.0
zip component.zip kinesis_data_stream.py requirements.txtaws s3 cp component.zip s3://$S3_BUCKET/artifacts/jp.co.xyz.StreamManagerKinesis/$VERSION/rm component.zipコンポーネントの登録
コンポーネントを登録するには、以下の手順を行います。
- 先ほど作成した
recipe.yamlファイルの内容をコピーします。 $ArtifactsS3Bucketを、コンポーネントのアーティファクトが実際に置かれているS3バケット名に置き換えます。

コンポーネントのデプロイ
Deployをクリックします。

Configure componentをクリックします。

必要な情報でコンポーネントの設定を更新します。KinesisStreamパラメータには、先ほど作成したKinesis Data Streamの名前を設定します。

設定内容を確認し、コンポーネントをデプロイします。

デプロイが完了すると、コンポーネントがGreengrassストリームから指定したKinesis Data Streamへデータをストリーミングするようになります。
テスト
S3バケットにオブジェクトが到着していれば、セットアップが正しく機能していることが確認できます。
aws s3 cp s3://<YOUR_BUCKET_NAME>/... ./以下の例は、読みやすさのために整形しています。
{ "propertyAlias": "/MyObject/MyVariable", "propertyValues": [ { "value": { "doubleValue": 7.699999999999997 }, "timestamp": { "timeInSeconds": 1661581962, "offsetInNanos": 9000000 }, "quality": "GOOD" } ]}データが期待通りに表示されれば、OPC UAデータがKinesis Data Streamに正しくストリーミングされ、S3に永続化されていることになります。
まとめ
ダミーのOPC UAサーバーをSiteWise Edge Gatewayとカスタムのgreengrassコンポーネント経由で橋渡ししたところ、工場現場を模したテレメトリをKinesis Data Streamに流し込み、S3に永続化できました。SiteWise Edge Gatewayは、サーバーをポーリングしてSiteWise_Stream_Kinesisのような名前付きGreengrassストリームに書き込むという、OPC UA側の処理を単独で完結させます。しかし、SiteWiseはKinesisへ直接エクスポートする機能を持たないため、そのデータをKinesisに取り込むには、ここで構築したカスタムコンポーネントが必要になります。コンポーネント自体はかなり小さく、Stream Manager SDKのStreamManagerClientとKinesisConfigが実際の転送処理を担い、BatchSize(上限500)が各PutRecords呼び出し前にメッセージをいくつ蓄積するかを制御します。SiteWiseが提供する機能と、このパターンが必要とする機能とのこのギャップは、あらかじめ計画しておく価値があります。つまり、コンポーネントのIAMロール、パッケージング、バージョン管理は、一度きりのセットアップ作業ではなく、パイプラインの継続的な一部として自ら所有し続けることになるからです。
Related posts
Docker環境でGreengrassコンポーネントをビルド・デプロイする
Greengrass Core Dockerイメージを使って、AWS IoT Greengrassコンポーネントをローカルで開発します。
AWS IoT Core における Kinesis Firehose のレコードセパレーター設定
IoT Core のトピックルールにある Firehose アクションでレコードセパレーターを設定し、S3 に改行区切りでレコードを格納する方法。
AWS IoT トピックペイロードによる Kinesis シャードの分離
IoT トピックのペイロードを顧客 ID ごとに異なる Kinesis シャードへルーティングし、後続処理でも顧客ごとの順序を保つ方法。
Kinesis Firehoseの動的パーティショニングでNDJSON処理を簡素化する
Kinesis Data Firehoseが動的パーティショニングをサポートしたことで、NDJSON(改行区切りJSON)への変換のためだけにLambda関数を用意する必要がなくなりました。
AthenaとKinesis Data FirehoseでS3のログをクエリする
Kinesis Data FirehoseでログをS3に流し込み、Athenaでクエリする。カスタムプレフィックスを使ってパーティションを自動登録させる方法。
