SiteWise Edge Gateway で OPC UA データを Kinesis にストリーミングする

SiteWise Edge Gateway で OPC UA データを Kinesis にストリーミングする

SiteWise Edge Gateway とカスタムの Greengrass コンポーネントを使い、OPC UA テレメトリを Kinesis Data Streams へ橋渡しします。

Takahiro Iwasa
7 min read

SiteWise Edge Gateway は、OPC UA のデータソースを Kinesis Data Streams に橋渡しし、工場のテレメトリを下流の AWS サービスで利用できるようにします。

Overview Diagram

バックエンドの構築

OPC UA サーバー

サンプルの OPC UA サーバーには、Python ライブラリの opcua-asyncio を使用します。

必要なパッケージをインストールします。

Terminal window
pip install asyncua

サンプルスクリプト(server-minimal.py)を EC2 インスタンスにダウンロードします。

Terminal window
curl -OL https://raw.githubusercontent.com/FreeOpcUa/opcua-asyncio/master/examples/server-minimal.py

スクリプトを実行します。

Terminal window
python server-minimal.py

Greengrass 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 を選択します。

Step 1

Data processing pack は不要です。

Step 2

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

Step 3

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

Step 4

Step 4

Step 4

設定内容を確認し、Create をクリックします。

Step 5

Kinesis Data Stream

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

Kinesis Data Stream Setup

Kinesis Data Firehose

ストリーミングされたデータを永続化するため、以下の内容で Kinesis Data Firehose 配信ストリームを設定します。

  • ソース:上記で作成した Kinesis Data Stream
  • 送信先:S3 バケット

Kinesis Firehose Setup

Greengrass コンポーネントの作成

Greengrass ストリームから Kinesis Data Streams へデータを転送するため、カスタムの Greengrass コンポーネントを作成します。コンポーネント開発の詳細は、公式ドキュメントを参照してください。

ディレクトリ構成

コンポーネントのファイルは、以下のような構成にします。

Terminal window
- kinesis_data_stream.py
- recipe.yaml
- requirements.txt
- stream_manager_sdk.zip

コンポーネントの作成

ℹ️ Note

stream_manager_sdk.zip には、Greengrass コンポーネントで使う SDK を含めます。詳細とサンプルコードは、AWS Greengrass Stream Manager SDK for Pythonを参照してください。

requirements.txt を作成します。stream-manager ライブラリを使って Greengrass Stream Manager にアクセスします。

requirements.txt
cbor2~=5.4.2
stream-manager==1.1.1

以下の内容で kinesis_data_stream.py を作成します。このスクリプトは、Greengrass Stream Manager SDK を使って Kinesis Data Stream へデータを転送します。

kinesis_data_stream.py
"""
Script to use Greengrass Stream Manager to stream data to a Kinesis Data Stream
See also https://github.com/aws-greengrass/aws-greengrass-stream-manager-sdk-python/blob/main/samples/stream_manager_kinesis.py
"""
import argparse
import asyncio
import logging
import 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)
recipe.yaml
# Replace $ArtifactsS3Bucket with your value to complete component registration.
RecipeFormatVersion: 2020-01-25
ComponentName: jp.co.xyz.StreamManagerKinesis
ComponentVersion: 1.0.0
ComponentDescription: Streams data in Greengrass stream to a Kinesis Data Stream.
ComponentPublisher: self
ComponentDependencies:
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 バケットにアップロードします。

Terminal window
S3_BUCKET=<YOUR_BUCKET_NAME>
VERSION=1.0.0
zip component.zip kinesis_data_stream.py requirements.txt
aws s3 cp component.zip s3://$S3_BUCKET/artifacts/jp.co.xyz.StreamManagerKinesis/$VERSION/
rm component.zip

コンポーネントの登録

コンポーネントを登録するには、以下の手順を行います。

  1. 先ほど作成した recipe.yaml の内容をコピーします。
  2. $ArtifactsS3Bucket を、コンポーネントのアーティファクトを置いた S3 バケット名に置き換えます。

Component Registration

コンポーネントのデプロイ

Deploy をクリックします。

Deploy

Configure component をクリックします。

Configure Component

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

Update Configuration

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

Deploy Component

デプロイが完了すると、コンポーネントが Greengrass ストリームから指定した Kinesis Data Stream へデータを転送します。

テスト

S3 バケットにオブジェクトが保存されていれば、セットアップが正しく機能しています。

Terminal window
aws s3 cp s3://<YOUR_BUCKET_NAME>/... ./
ℹ️ Note

以下の例は、読みやすさのために整形しています。

{
"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 の StreamManagerClientKinesisConfig が転送を担い、BatchSize(上限 500)で各 PutRecords 呼び出しにまとめるメッセージ数を制御します。

コンポーネントの IAM ロール、パッケージ、バージョンは、一度きりの設定ではなく、パイプラインの一部として継続的に管理する必要があります。

About the author

Takahiro Iwasa

Takahiro Iwasa

Software Developer

This blog shares technical notes from hands-on projects—architecture, implementation, and AWS service integrations.