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

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

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

Takahiro Iwasa
8 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コンポーネントを作成する必要があります。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を作成します。Greengrass Stream Managerとやり取りするには、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に永続化されていることになります。

まとめ

ダミーの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のStreamManagerClientKinesisConfigが実際の転送処理を担い、BatchSize(上限500)が各PutRecords呼び出し前にメッセージをいくつ蓄積するかを制御します。SiteWiseが提供する機能と、このパターンが必要とする機能とのこのギャップは、あらかじめ計画しておく価値があります。つまり、コンポーネントの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.