第4章 Camel Kafka コネクターの拡張
本章では、Camel Kafka コネクターおよびコンポーネントを拡張し、カスタマイズする方法を説明します。Camel Kafka Connector は、コードを作成せずに Kafka Connect フレームワークで Camel コンポーネントを直接設定する簡単な方法を提供します。しかし、状況によっては、特定のユースケースに対して Camel Kafka Connector を拡張およびカスタマイズする場合があります。
4.1. Camel Kafka コネクターアグリゲーターの設定 リンクのコピーリンクがクリップボードにコピーされました!
Camel Kafka シンクコネクターを使用する一部のシナリオでは、Kafka レコードを外部シンクシステムに送信する前に、アグリゲーターを追加して Kafka レコードをバッチ処理することがあります。通常、このような場合、レコードの集約のために特定のバッチサイズとタイムアウトを定義します。完了したら、集約レコードが外部システムに送信されます。
Apache Camel によって提供されるアグリゲーターの 1 つを使用して、Camel Kafka Connector プロパティーで集約の設定を指定できます。または、Java でカスタムアグリゲーターを実装することもできます。ここでは、Camel Kafka Connector プロパティーで Camel アグリゲーターを設定する方法を説明します。
前提条件
- Camel Kafka Connector がインストール済みである必要があります (例: 「OpenShift での AMQ Streams および Kafka Connect S2I のインストール」 を参照)。
- シンクコネクターがデプロイ済みである必要があります (例: 「OpenShift での Kafka Connect S2I を使用した Camel Kafka コネクターのデプロイ」 を参照)。ここでは、AWS S3 シンクコネクターを使用する例を紹介します。
手順
インストールプラットフォームに応じて、シンクコネクターおよびアグリゲーターを Camel Kafka Connector プロパティーで設定します。
- OpenShift
以下の例は、カスタムリソースの AWS S3 シンクコネクターおよびアグリゲーター設定を示しています。
oc apply -f - << EOF apiVersion: kafka.strimzi.io/v1alpha1 kind: KafkaConnector metadata: name: s3-sink-connector namespace: myproject labels: strimzi.io/cluster: my-connect-cluster spec: class: org.apache.camel.kafkaconnector.aws2s3.CamelAws2s3SinkConnector tasksMax: 1 config: key.converter: org.apache.kafka.connect.storage.StringConverter value.converter: org.apache.kafka.connect.storage.StringConverter topics: s3-topic camel.sink.path.bucketNameOrArn: camel-kafka-connector camel.sink.endpoint.keyName: ${date:now:yyyyMMdd-HHmmssSSS}-${exchangeId} # Camel aggregator settings camel.beans.aggregate: #class:org.apache.camel.kafkaconnector.aggregator.StringAggregator camel.beans.aggregation.size: 10 camel.beans.aggregation.timeout: 5000 camel.component.aws2-s3.accessKey: xxxx camel.component.aws2-s3.secretKey: yyyy camel.component.aws2-s3.region: region EOF- Red Hat Enterprise Linux
次の例は、
CamelAwss3SinkConnector.propertiesファイル内の AWS S3 シンクコネクターとアグリゲーターの設定を示しています。name=CamelAWS2S3SinkConnector connector.class=org.apache.camel.kafkaconnector.aws2s3.CamelAws2s3SinkConnector key.converter=org.apache.kafka.connect.storage.StringConverter value.converter=org.apache.kafka.connect.storage.StringConverter topics=mytopic camel.sink.path.bucketNameOrArn=camel-kafka-connector camel.component.aws2-s3.access-key=xxxx camel.component.aws2-s3.secret-key=yyyy camel.component.aws2-s3.region=eu-west-1 camel.sink.endpoint.keyName=${date:now:yyyyMMdd-HHmmssSSS}-${exchangeId} # Camel aggregator settings camel.beans.aggregate=#class:org.apache.camel.kafkaconnector.aggregator.StringAggregator camel.beans.aggregation.size=10 camel.beans.aggregation.timeout=5000