Apache CamelフレームワークでのSAS Event Stream Processingの使用
- 概要
- Apache Camel Frameworkのインストール
- RabbitMQライブラリのインストール
- SAS Event Stream Processingの実装
- MavenプロジェクトでCamelコンポーネントを使用
- エンドポイントの設定
- 変換Beanの使用
- 例
概要
Apache Camelフレームワークを使用すると、さまざまなアプリケーションを1つのまとまったアーキテクチャに統合できます。Apache Camelフレームワークを使用して、エンドポイントを含むルートを設定します。
<route id="injectTrades" startupOrder="10">
<from uri="systemA://someThing"/>
<to uri="systemB://someOtherThing"/>
</route>
このコードの fromエレメントと toエレメントは、コンポーネントを参照するApache Camelエンドポイントです。Apache Camelでは、すぐに使用できる多くのコンポーネントがサポートされています。カスタムコンポーネントを開発するためのツールが提供されます。Apache
Camel Consumerはfromエンドポイントにマップし、Apache Camel Producerはtoエンドポイントにマップします。
Apache Camelフレームワークの詳細については、documentationを参照してください。
Apache Camel Frameworkのインストール
SAS Event Stream ProcessingでApache Camelフレームワークを使用するには、さまざまなファイルをダウンロードしてインストールする必要があります。これには、Apacheコンポーネントと特定のJARファイルが含まれます。
- 次のApacheコンポーネントにアクセスしてインストールします。
- Apache Camel(http://camel.apache.org/からダウンロードできます)
- Apache Maven(https://maven.apache.org/からダウンロードできます)
Apache Mavenは、SAS Event Stream Processingのコンポーネントを活用するプロジェクトの作成に使用できるビルド環境です。Apache Mavenをインストールするときは、パスにMavenインストールディレクトリの下位のbinディレクトリをインストールしてください。
- Apacheコンポーネントをインストールしたら、2つのJARファイルをローカルのMavenリポジトリにインストールします。
- ESP APIクライアントJAR
- ESP Camel JAR
注: これらのJARファイルは、$DFESP_HOME/libにあります。それらにはクライアントAPIとCamelコンポーネントが含まれています。 - JARファイルをMavenリポジトリにインストールします。
$ cd $DFESP_HOME/lib $ mvn install:install-file -Dfile=dfx-esp-api.jar -DgroupId=com.sas.esp -DartifactId=dfx-esp-api -Dversion=6.1 -Dpackaging=jar $ mvn install:install-file -Dfile=dfx-esp-camel.jar -DgroupId=com.sas.esp -DartifactId=dfx-esp-camel -Dversion=6.1 -Dpackaging=jar $ mvn install:install-file -Dfile=cas-client-3.06.jar -DgroupID=com.sas.esp -DartifactID=dfx-cas-auth -Dversion=3.0.6 -Dpackaging=jar - JARファイルをインストールしたら、Mavenプロジェクトオブジェクトモデル(pom.xml)ファイルからそれらを参照できます。イベントストリームプロセッサクライアントAPIとSAS Event Stream Processing Camelコンポーネントの両方のエントリが必要です。
イベントストリームプロセッサクライアントAPIエントリは次のとおりです。
<dependency>
<groupId>com.sas.esp</groupId>
<artifactId>dfx-esp-api</artifactId>
<version>[6.1]</version>
</dependency>
SAS Event Stream Processing Camelのエントリは次のとおりです。
<dependency>
<groupId>com.sas.esp</groupId>
<artifactId>dfx-esp-camel</artifactId>
<version>[6.1]</version>
</dependency>
CAS認証エントリは次のとおりです。
<dependency>
<groupID>com.sas.esp</groupID>
<artifactID>dfs-cas-auth</artifactID>
<version>3.0.6</version>
</dependency>
RabbitMQライブラリのインストール
SAS Event Stream ProcessingがRabbitMQを代替トランスポートライブラリとして使用するようにMavenプロジェクトを設定します。
- RabbitMQ API JARをインストールします。
$ cd $DFESP_HOME/lib $ mvn install:install-file -Dfile=dfx-esp-rabbitmq-api.jar -DgroupId=com.sas.esp -DartifactId=dfx-esp-rabbitmq-api -Dversion=6.1 -Dpackaging=jar - RabbitMQ依存関係情報を含むようにMavenプロジェクトオブジェクトモデル(pom.xml)を更新します。
<dependency> <groupId>com.rabbitmq</groupId> <artifactId>amqp-client</artifactId> <version>3.5.6</version> </dependency> <dependency> <groupId>commons-configuration</groupId> <artifactId>commons-configuration</artifactId> <version>1.10</version> </dependency> <dependency> <groupId>com.sas.esp</groupId> <artifactId>dfx-esp-rabbitmq-api</artifactId> <version>6.1</version> </dependency>
dfx-esp-api)の依存関係を定義する前に、pom.xmlファイル内のRabbitMQ JAR (dfx-esp-rabbitmq-api)の依存関係を定義する必要があります。SAS Event Stream Processingの実装
SAS Event Stream Processingの実装は、コンシューマー(パブリッシュ/サブスクライブサブスクライバーを実装する)またはプロデューサー(パブリッシュ/サブスクライブパブリッシャーを実装する)のいずれかであるSAS
Event Stream Processing Camelエンドポイントで構成されます。コンシューマーはfromエンドポイントにマップし、プロデューサーはtoエンドポイントにマッピングします。
たとえば、1つのパブリッシュ/サブスクライブサーバーからイベントを受信して別のパブリッシュ/サブスクライブサーバーに送信するには、次のコードを実行します。
...
<endpoint id="subscribe" uri="esp://espsrv01:46003">
<property key="project" value="project" />
<property key="contquery" value="query" />
<property key="window" value="transform" />
</endpoint>
<endpoint id="publish" uri="esp://espsrv01:47003">
<property key="project" value="project" />
<property key="contquery" value="query" />
<property key="window" value="trades" />
</endpoint>
<route>
<from uri="ref:subscribe"/>
<to uri="ref:publish" />
</route>
...
SAS Event Stream Processingコンポーネントは、イベントを表す次の形式で動作します。
|
形式 |
説明 |
|---|---|
|
マップ<String,Object> |
キーがフィールド名であり、値がフィールド値である標準Java Mapオブジェクトです。 |
|
リスト<Map<String,Object>> |
キーがフィールド名であり、値がフィールド値であるJava Mapオブジェクトの標準JAVAリスト。 |
|
XML |
XML形式のイベントデータ。 |
|
JSON |
JSONのイベントデータ。 |
SAS Event Stream Processingコンシューマー(サブスクライバー)はイベントを受信し、これらの形式の1つに変換してルートに沿って送信します。SAS Event Stream Processing Producer(パブリッシャー)は、これらの形式のいずれかでデータを受信し、イベントに変換してパブリッシュします。
これらの標準形式により、SAS Event Stream Processingは、Camelフレームワークで利用可能な他のコンポーネントと容易に統合され、それらとデータを共有します。データを必要な形式にするために何らかの変換が必要な場合は、変換Beanをエンドポイント間で使用できます。
最初にcsvDataというエンドポイントを作成し、ファイルからcsvデータを読み取ります。また、ソースウィンドウにこのデータをインジェクトするためのinjectというエンドポイントを作成します。使用するルートはファイルからSAS Event Stream Processingになりますが、サポートされている形式の1つにcsvデータを取得する必要があります。
CSVデータを適切な形式に変換するには、スキーマと出力形式を必要とするcsvTransformという名前のビーンを作成します。その後、csvデータからイベントを作成し、ルートを通るXMLドキュメントを作成できます。これはプロデューサーによって消費され、プロデューサーは指定されたソースウィンドウにデータをインジェクトします。いずれかのSAS
Event Stream Processing変換Beanのメソッド属性は常に変換です。
csvTransformへの入力は String型でなければなりません。CSVデータをPOST/HTTPリクエストのボディとして変換ビーンに渡すと、InputStreamとして渡されるため、エラーが発生します。Camel
が提供しているconvertBodyToメソッドを使って、HTTPボディを文字列に変換することができます。カンマで区切られた標準の値のイベント形式をサポートされている型の1つに変換する例を次に示します。
...
<endpoint id="csvData" uri="stream:file">
<property key="fileName" value="/mnt/data/share/tradesData/trades1M.csv" />
</endpoint>
<endpoint id="inject" uri="esp://espsrv01:46003">
<property key="project" value="project" />
<property key="contquery" value="query" />
<property key="window" value="trades" />
</endpoint>
<route id="injectTrades" startupOrder="10">
<from uri="ref:csvData"/>
<bean ref="csvTransform" method="transform" />
<to uri="ref:inject"/>
</route>
<bean id="csvTransform" class="com.sas.esp.camel.transforms.CsvTransform">
<property name="schema" value="id*:int64,symbol:string,currency:
int32,time:int64,msecs:int32,price:double,quant:int32,venue:int32,broker:
int32,buyer:int32,seller:int32,buysellflg:int32" />
<property name="format" value="xml" />
</bean>
...
MavenプロジェクトでCamelコンポーネントを使用
SAS Event Stream Processing CamelコンポーネントをURIで参照するには、次のファイルを作成する必要があります。
META-INF/services/org/apache/camel/component/esp
このファイルには、次の行が含まれている必要があります。
class=com.sas.esp.clients.camel.EspComponent
これにより、次のようなURIを持つイベントストリーム処理コンポーネントを参照できます。
<endpoint id="inject" uri="esp://<pub/sub host>:<pub/sub port>">
エンドポイントの設定
通信しようとしているSAS Event Stream Processingパブリッシュ/サブスクライブサーバーのホストとポートを決定する必要があります。この情報を使用して、コンポーネントのURIを指定します。たとえば、パブリッシュ/サブスクライブサーバーがマシンespsrv01のポート46003で実行されている場合、URIは次のようになります。
esp://espsrv01:46003
常にイベントをウィンドウに入れたり、ウィンドウからイベントを取得したりするため、プロジェクト、連続クエリ、および関心のあるウィンドウも指定する必要があります。次のいずれかの方法を使用してこれを行うことができます。
- URIにパラメーターを追加する
esp://espsrv01:46003?project=myproject&contquery=mycq&window=trades - プロパティを持つエンドポイントエレメントを指定する
<endpoint id="inject" uri="esp://espsrv01:46003"> <property key="project" value="myproject" /> <property key="contquery" value="mycq" /> <property key="window" value="trades" /> </endpoint>
|
プロパティ |
説明 |
コンシューマー |
プロデューサー |
有効な値 |
|---|---|---|---|---|
|
|
イベントストリーム処理プロジェクト。 |
x |
x |
有効なプロジェクトです。 |
|
|
イベントストリーム処理の連続クエリ。 |
x |
x |
有効な連続クエリ。 |
|
|
イベントストリーム処理ウィンドウ。 |
x |
x |
有効なウィンドウ。 |
|
|
このコンポーネントのデータ形式。これは通常、経路を送信するイベントデータの形式を指定するためにコンシューマーとともに使用されます。これはプロデューサーが使用できます。コンポーネントが文字列である本文を含むメッセージを受け取った場合、プロデューサーは変換に形式(XMLまたはJSON)を使用します。 |
x |
x |
マップ、リスト、XML、JSON |
|
|
日付フィールドの出力日付形式。 |
x |
Java SimpleDateFormat | |
|
|
ソースウィンドウにインジェクトするイベントブロックのサイズ。 |
x |
符号なし整数 | |
|
|
メッセージ本文全体を格納するために使用されるイベントフィールドの名前。このプロパティが設定されている場合、プロデューサーはメッセージを受信し、キーがblobプロパティによって示されるフィールドであるMapMap<String,Object>を作成します。値は、文字列として表されるメッセージ本文全体です。このデータは、適切なソースウィンドウにインジェクトされます。 注: blobプロパティを使用する場合、イベントがインジェクトされるソースウィンドウは、 'insert-only = true'と 'autogen-key = true'の両方を設定する必要があります。 |
x |
有効なイベントフィールド | |
|
|
SASLogon認証で使用する名前。 |
x |
x | |
|
|
OAuth認証で使用するOAuthトークン。 |
x |
x |
変換Beanの使用
データがイベントストリーム処理コンポーネントによってすぐに使用可能な形式でない場合は、変換Beanを使用してデータを使用可能な形式に変換する必要があります。変換Beanを使用して、SAS Event Stream Processing csv形式のイベントを使用可能な形式に変換できます。変換Beanは、fromとendの間のルートに配置します。
次に例を示します。
...
<route id="injectTrades">
<from uri="ref:csvData"/>
<bean ref="csvTransform" method="transform" />
<to uri="ref:inject"/>
</route>
...
<bean id="csvTransform" class="com.sas.esp.clients.camel.transforms.CsvTransform">
<property name="schema" value="id*:int64,symbol:string,currency:
int32,time:int64,msecs:int32,price:double,quant:int32,venue:int32,broker:
int32,buyer:int32,seller:int32,buysellflg:int32" />
<property name="format" value="xml" />
<property name="dateformat" value="ddMMMyyyy HH:mm:ss.S a" />
</bean>
...
各Beanは、SAS Event Stream Processingで使用できるようにデータを変換するのに役立つ特定のパラメーターを取ります。csv変換Beanには、イベントスキーマと出力データ形式が必要です。Beanは、ファイル参照を含むfromエンドポイントとイベントをウィンドウにパブリッシュするtoエンドポイントの間にあります。
|
クラス |
説明 |
|---|---|
|
com.sas.esp.clients.camel.transforms.CsvTransform |
csvデータをESPイベントデータに変換する
|
|
com.sas.esp.clients.camel.transforms.RssTransform |
RSSデータをESPイベントデータに変換する
|
例
サンプルの入手先
サンプルはサポートWebサイトからダウンロードできます。java/camelに移動し、 csv 、distributed 、およびrssサブディレクトリを調べます。
CSVインジェクション
次の例では、csvファイルからトレードデータを読み取り、その取引をブローカ監視モデルにインジェクトします。また、brokerAlertsAggrウィンドウにサブスクライブし、これらのイベントをJSON形式でコンソールに書き込みます。
- src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
espServer=esp://espsrv01:46003 tradesFile=data/trades1M.csv - ESPサーバーを起動します。
$ dfesp_xml_server -model file://model.xml -http-admin <http admin port> -http-pubsub <http pub/sub port> -pubsub <esp pub/sub port> -nocleanup - プロジェクトを開始します。
$ mvn camel:run
分散モデリング
この例では、2つのサーバー間にブローカ監視モデルを配布します。第1のサーバは、トレードデータを受け取り、イベントへのすべてのディメンション追加(ブローカ情報、会場データ)を実行します。このモデルはprimary.xmlにあり、変換という機能ウィンドウで終わります。変換によって生成されるイベントには、一連のトレード情報が含まれます。プロジェクトはserver1のこのウィンドウにサブスクライブし、イベントをserver2に転送します。server2はブローカのアラートを探します。別のルートを使用して、server2のbrokerAlertsAggrウィンドウをサブスクライブし、イベントをマップ形式で画面にダンプします。
- src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
espServer1=esp://espsrv01:46003 espServer2=esp://espsrv01:47003 tradesFile=data/trades1M.csv - プライマリEvent Stream Processing Serverを起動します。
$ dfesp_xml_server -model file://primary.xml -http-admin <http admin port> -http-pubsub <http pub/sub port> -pubsub <esp pub/sub port> -nocleanup - セカンダリEvent Stream Processing Serverを起動します。
$ dfesp_xml_server -model file://secondary.xml -http-admin <http admin port> -http-pubsub <http pub/sub port> -pubsub <esp pub/sub port> -nocleanup - プロジェクトを開始します。
$ mvn camel:run
RSS
この例では、Camel RSSコンポーネントを使用して、任意の数のRSSフィードからデータを読み込んで、それらをSAS Event Stream Processingにインジェクトするルートを設定します。ルートにRSSフィードを追加できるはずです。
<route>
<from uri="rss:http://feeds.reuters.com/reuters/businessNews" />
<from uri="rss:http://feeds.reuters.com/reuters/topNews" />
<from uri="rss:http://feeds.reuters.com/reuters/technologyNews" />
<bean ref="rssTransform" method="transform" />
<to uri="ref:publishNews" />
</route>
...
新しい変換Beanを使用して、RSSデータがサポートされている形式にも変換できます。
<bean id="rssTransform" class="com.sas.esp.clients.camel.transforms.RssTransform">
<property name="opcode" value="upsert" />
</bean>
...
次のプロジェクトでは、RSSデータはタイトルによってキーイングされています。
<project name='project' pubsub='auto' threads='4'>
<contqueries>
<contquery name='cq' trace='src'>
<windows>
<window-source name='src'>
<schema-string>title*:string,author:string,link:string,
description:string,categories:string,pubDate:date</schema-string>
</window-source>
</windows>
</contquery>
</contqueries>
</project>
...
- src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
espServer=esp://espsrv01:46003 - Event Stream Processing Serverを起動します。
$ dfesp_xml_server -model file://model.xml -http-admin <http admin port> -http-pubsub <http pub/sub port> -pubsub <esp pub/sub port> -nocleanup - プロジェクトを開始します。
$ mvn camel:run
天気
この例では、Camel Weather Componentを使用して、任意の数の場所の気象データを読み取り、SAS Event Stream Processingにインジェクトするルートを設定します。任意の場所をルートに追加できます。次の例に示すように、ロケーションをエンドポイントとして定義できます。
<endpoint id="cary" uri="weather:foo">
<property key="location" value="cary,nc"/>
<property key="mode" value="XML"/>
<property key="units" value="IMPERIAL"/>
</endpoint>
<endpoint id="morehead" uri="weather:foo">
<property key="location" value="moreheadcity,nc"/>
<property key="mode" value="XML"/>
<property key="units" value="IMPERIAL"/>
</endpoint>
<endpoint id="chapelHill" uri="weather:foo">
<property key="location" value="chapelhill,nc"/>
<property key="mode" value="XML"/>
<property key="units" value="IMPERIAL"/>
</endpoint>
...
これらのエンドポイントをルートに追加できます。
<route>
<from uri="ref:cary"/>
<from uri="ref:morehead"/>
<from uri="ref:chapelHill"/>
<to uri="ref:publishWeather"/>
</route>
...
- src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
espServer=esp://espsrv01:46003 - Event Stream Processing Serverを起動します。
$ dfesp_xml_server -model file://model.xml -http-admin <http admin port> -http-pubsub <http pub/sub port> -pubsub <esp pub/sub port> -nocleanup - プロジェクトを開始します。
$ mvn camel:run