Apache CamelフレームワークでのSAS Event Stream Processingの使用

概要

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ファイルが含まれます。

  1. 次のApacheコンポーネントにアクセスしてインストールします。
    • Apache Camel(http://camel.apache.org/からダウンロードできます)
    • Apache Maven(https://maven.apache.org/からダウンロードできます)

    Apache Mavenは、SAS Event Stream Processingのコンポーネントを活用するプロジェクトの作成に使用できるビルド環境です。Apache Mavenをインストールするときは、パスにMavenインストールディレクトリの下位のbinディレクトリをインストールしてください。

  2. Apacheコンポーネントをインストールしたら、2つのJARファイルをローカルのMavenリポジトリにインストールします。
    • ESP APIクライアントJAR
    • ESP Camel JAR
    注: これらのJARファイルは、$DFESP_HOME/libにあります。それらにはクライアントAPIとCamelコンポーネントが含まれています。
  3. 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
  4. 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プロジェクトを設定します。

  1. 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
    
  2. 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>
注: ESP APIクライアントJAR (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
注: これについては、http://camel.apache.org/writing-components.htmlで詳しく説明しています。

これにより、次のような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>
エンドポイントで設定できるプロパティ

プロパティ

説明

コンシューマー

プロデューサー

有効な値

project

イベントストリーム処理プロジェクト。

x

x

有効なプロジェクトです。

contquery

イベントストリーム処理の連続クエリ。

x

x

有効な連続クエリ。

window

イベントストリーム処理ウィンドウ。

x

x

有効なウィンドウ。

format

このコンポーネントのデータ形式。これは通常、経路を送信するイベントデータの形式を指定するためにコンシューマーとともに使用されます。これはプロデューサーが使用できます。コンポーネントが文字列である本文を含むメッセージを受け取った場合、プロデューサーは変換に形式(XMLまたはJSON)を使用します。

x

x

マップ、リスト、XML、JSON

dateformat

日付フィールドの出力日付形式。

x

Java SimpleDateFormat

blocksize

ソースウィンドウにインジェクトするイベントブロックのサイズ。

x

符号なし整数

blob

メッセージ本文全体を格納するために使用されるイベントフィールドの名前。このプロパティが設定されている場合、プロデューサーはメッセージを受信し、キーがblobプロパティによって示されるフィールドであるMapMap<String,Object>を作成します。値は、文字列として表されるメッセージ本文全体です。このデータは、適切なソースウィンドウにインジェクトされます。

注: blobプロパティを使用する場合、イベントがインジェクトされるソースウィンドウは、 'insert-only = true'と 'autogen-key = true'の両方を設定する必要があります。

x

有効なイベントフィールド

authUser

SASLogon認証で使用する名前。

x

x

oAuthToken

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エンドポイントの間にあります。

SAS Event Stream Processing Camelパッケージの変換Bean

クラス

説明

com.sas.esp.clients.camel.transforms.CsvTransform

csvデータをESPイベントデータに変換する

schema

csvデータによって表されるイベントのイベントスキーマ。

形式

イベントの書き込みに使用するデータ形式。有効な値は、map、list、XML、およびJSONです。

dateformat

使用する日付形式。

com.sas.esp.clients.camel.transforms.RssTransform

RSSデータをESPイベントデータに変換する

形式

イベントの書き込みに使用するデータ形式。有効な値は、list、XML、およびJSONです。

例

サンプルの入手先

サンプルはサポートWebサイトからダウンロードできます。java/camelに移動し、 csv 、distributed 、およびrssサブディレクトリを調べます。

CSVインジェクション

次の例では、csvファイルからトレードデータを読み取り、その取引をブローカ監視モデルにインジェクトします。また、brokerAlertsAggrウィンドウにサブスクライブし、これらのイベントをJSON形式でコンソールに書き込みます。

  1. src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
    espServer=esp://espsrv01:46003
    tradesFile=data/trades1M.csv
    
  2. 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
  3. プロジェクトを開始します。
    $ mvn camel:run
    

分散モデリング

この例では、2つのサーバー間にブローカ監視モデルを配布します。第1のサーバは、トレードデータを受け取り、イベントへのすべてのディメンション追加(ブローカ情報、会場データ)を実行します。このモデルはprimary.xmlにあり、変換という機能ウィンドウで終わります。変換によって生成されるイベントには、一連のトレード情報が含まれます。プロジェクトはserver1のこのウィンドウにサブスクライブし、イベントをserver2に転送します。server2はブローカのアラートを探します。別のルートを使用して、server2のbrokerAlertsAggrウィンドウをサブスクライブし、イベントをマップ形式で画面にダンプします。

  1. src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
    espServer1=esp://espsrv01:46003
    espServer2=esp://espsrv01:47003
    tradesFile=data/trades1M.csv
  2. プライマリ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
  3. セカンダリ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
  4. プロジェクトを開始します。
    $ 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>

...
  1. src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
    espServer=esp://espsrv01:46003
  2. 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
  3. プロジェクトを開始します。
    $ 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>

...
  1. src/main/resources/esp.propertiesを編集して、パブリッシュ/サブスクライブサーバー情報が含まれるようにします。
    espServer=esp://espsrv01:46003
  2. 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
  3. プロジェクトを開始します。
    $ mvn camel:run
    
最終更新: 2025年4月23日