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 コンポーネントを作成します。コンポーネント開発の詳細は、公式ドキュメントを参照してください。
ディレクトリ構成
コンポーネントのファイルは、以下のような構成にします。
- 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 を作成します。stream-manager ライブラリを使って Greengrass 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 に永続化されています。
まとめ
SiteWise Edge Gateway とカスタムの Greengrass コンポーネントを使い、サンプルの OPC UA サーバーから Kinesis Data Streams へテレメトリを転送して、S3 に永続化できました。
SiteWise Edge Gateway は OPC UA サーバーをポーリングし、SiteWise_Stream_Kinesis のような名前付き Greengrass ストリームへデータを書き込みます。SiteWise から Kinesis へ直接エクスポートすることはできないため、カスタムコンポーネントでメッセージを転送します。Stream Manager SDK の StreamManagerClient と KinesisConfig が転送を担い、BatchSize(上限 500)で各 PutRecords 呼び出しにまとめるメッセージ数を制御します。
コンポーネントの 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 の動的パーティショニングを使い、変換用の Lambda 関数を用意せずに NDJSON(改行区切り JSON)を S3 へ出力します。
Athena と Kinesis Data Firehose で S3 のログをクエリする
Kinesis Data Firehose でログを S3 に配信し、Athena でクエリします。カスタムプレフィックスを使ってパーティションを自動登録する方法を解説します。
