Apache Pulsar のカタログメタデータを Flink で取得する方法

Apache Pulsar に格納されているネームスペースやトピックのメタデータを、pulsar-flink-connector を用いて Java から取得する方法を紹介します。 Maven 依存関係 <dependency> <groupId>io.streamnative.connectors</groupId> <artifactId>pulsar-flink-connector-2.11-1.12</artifactId> <version>2.7.3</version> &l ...

7月7日 00:44 投稿

Flink と Kafka のオフセット管理方法

Flink と Kafka を連携する際、オフセット管理は重要な課題です。 自動管理モードでは、以下のような問題が発生します: プロセス途中での停止によりデータが失われることがあります 再起動時、同じデータが再び処理される可能性があります これらの問題を解決するためには、Kafka のオフセットを手動で管理し、Flinkのチェックポイントとオフセットを同期させる必要があ ...

6月21日 22:25 投稿

Flink学習メモ:ストリーム処理の基本とKafka統合

Flinkは「ストリーム処理」を基盤とした分散処理エンジンで、データが到着次第すぐに処理を行います。これはバッチ処理(例:Spark)とは異なり、マルチスレッドで逐次的に動作します。以下に、テキストファイルから読み込んだ単語をカウントする最小限の例を示します。 import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apach ...

6月1日 23:33 投稿