Info2softは、ウェブサイトでより快適で適切な閲覧体験を提供するためにCookieを使用しています。 プライバシーポリシー
Loading...
現代のイベント駆動型アーキテクチャでは、KafkaからPostgreSQLへのデータストリーミングは一般的な要件です。Kafkaはスケーラブルなリアルタイムイベント処理を実現し、PostgreSQLはアプリケーション、分析、データ同期向けに信頼性の高いリレーショナルストレージを提供します。
本ガイドでは、KafkaのデータがPostgreSQLへ流れる仕組み、一般的な連携手法、信頼性の高いリアルタイムデータパイプラインを構築するための主要なベストプラクティスを解説します。
リアルタイムデータパイプラインにおいて、Apache Kafkaはイベントストリーミングプラットフォームとして動作し、PostgreSQLは処理済みデータを保存・クエリするリレーショナルデータベースとして機能します。
KafkaからPostgreSQLへのデータストリーミングとは、Kafkaトピックからイベントを継続的に消費し、低レイテンシでPostgreSQLテーブルへ書き込むことを指します。
この連携は主に以下の用途で利用されます。
Kafkaエコシステムにおいて、ソースコネクタはデータをKafkaへ送り込み、シンクコネクタはKafkaからPostgreSQLなど外部システムへデータを転送します。
Kafkaからのデータストリーミングは、データベース同期や高可用性シナリオでよく使われるCDCやデータベースレプリケーションとは異なります。
| 連携パターン | Kafka→PG ストリーミング | CDC | DBレプリケーション |
|---|---|---|---|
| 目的 |
Kafkaイベントデータを PostgreSQLテーブルに書き込む |
データベースの変更をKafka または他システムへストリームする |
環境間でデータベースのコピー を維持する |
| データ単位 |
イベントメッセージ(JSON、Avro、 Protobuf) |
INSERT/UPDATE/DELETEイベント | トランザクションまたはログレコード |
| 主なユースケース |
リアルタイム分析、 データ統合 |
リアルタイム同期、イベント駆動アーキテクチャ | 災害復旧、高可用性、読み取りスケーリング |
信頼性の高いKafka‑to‑PostgreSQLパイプラインは、Kafkaトピックからリレーショナルテーブルへデータを移送する複数のコンポーネントで構成されます。このデータフローを理解することで、Kafka Connectやカスタムコンシューマーといった各連携手法の裏側の動作を把握できます。
Kafkaプロデューサーは、Kafkaトピックへレコードを発行するアプリケーションです。これらのレコードはキー‑バリュー形式のメッセージとして保存され、通常JSON、Apache Avro、Protocol Buffersなどの形式でシリアライズされます。
Kafkaトピックはパーティションに分割され、スケーラビリティと並列処理を実現します。プロデューサーがメッセージを送信すると、Kafkaはメッセージキーに基づいてパーティションに割り当てます。同一キーのメッセージは同一パーティションに書き込まれ、処理順序が保たれます。
受信側では、コンシューマーまたは連携コネクタがKafkaトピックからメッセージを読み取り、PostgreSQLテーブルへ書き込みます。代表的な処理フローは以下の通りです。
PostgreSQLへの書き込み処理はKafkaからの読み取りよりもレイテンシが高くなる傾向があるため、スループットを維持しオーバーヘッドを抑えるには、複数レコードをバッチ処理してデータベースへ書き込むことが重要です。
KafkaからPostgreSQLへデータをストリーミングする一般的な手法は3種類あります。適切な選択は、データ量、変換要件、レイテンシの期待値、運用の複雑さなどの要因に依存します。
Kafka ConnectはApache Kafkaエコシステムのオープンソースフレームワークで、Kafkaと外部システム間のデータ移送を実現します。
JDBCシンクコネクタは設定ベースのコネクタで、カスタムコンシューマーコードを作成せずにKafkaトピックからレコードを消費し、PostgreSQLテーブルへ書き込みます。
JDBCシンクコネクタの動作
JDBCシンクコネクタはKafka Connectクラスター内で動作し、設定されたKafkaトピックからレコードを継続的に消費します。コネクタ設定に基づいてKafkaレコードをデータベース操作に変換し、JDBCを介してPostgreSQLへデータを書き込みます。
コネクタはKafkaコンバーターを通じJSON、Avroなど対応フォーマットでシリアライズされたデータを処理可能です。Avroのようなスキーマベースのフォーマットでは、スキーマレジストリと連携しスキーマ互換性を管理するのが一般的です。
コネクタ設定例
下記はKafkaトピックの注文データをPostgreSQLへ書き込む代表的なJDBCシンクコネクタの設定例です。
{
"name": "postgres-order-sink",
"config": {
"connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
"tasks.max": "2",
"topics": "customer_orders",
"connection.url": "jdbc:postgresql://postgres-db.internal:5432/ecommerce",
"connection.user": "app_user",
"connection.password": "secure_db_password_123",
"insert.mode": "upsert",
"pk.mode": "record_key",
"pk.fields": "order_id",
"auto.create": "true",
"auto.evolve": "true"
}
}
auto.evolveは開発・テストを簡略化できますが、本番環境では慎重に使用してください。上流のデータ構造が適切なレビューなしに変更された場合、予期せぬデータベース変更が発生する可能性があります。適したユースケース
JDBCシンクコネクタによるKafka Connectは、スキーマが比較的安定した単純なデータ取り込みパイプラインに適しています。カスタムコンシューマーの作成・保守をせず、Kafka‑to‑PostgreSQLのデータ移送を安定的に実現したいチームに向いています。
データ変換、フィルタリング、ルーティング、アプリケーション固有のロジックを完全に制御したい場合、カスタムコンシューマーは柔軟な選択肢となります。既製のコネクタを使用せず、PythonやGoなどの言語で軽量アプリケーションを作成し、Kafkaメッセージを消費してPostgreSQLへ書き込みます。
(Kafkaトピック)—>[カスタムコンシューマーアプリケーション(Python/Go)]—>[PostgreSQLデータベース]
Pythonによるデータ消費と書き込み
カスタムコンシューマーはKafkaクラスターに直接接続し、新規レコードを取得し、アプリケーション層でメッセージを処理した後、結果をPostgreSQLへ書き込みます。
下記の例はconfluent‑kafka PythonクライアントでJSONメッセージを消費し、psycopg2を使用してupsert処理でPostgreSQLへレコードを書き込むサンプルです。
import json
import psycopg2
from confluent_kafka import Consumer, KafkaError
# Kafka consumer configuration
kafka_config = {
'bootstrap.servers': 'kafka.internal:9092',
'group.id': 'postgres-ingest-group',
'auto.offset.reset': 'earliest',
'enable.auto.commit': False
}
# PostgreSQL database connection
db_conn = psycopg2.connect("host=postgres-db.internal dbname=ecommerce user=app_user password=secure_password")
db_cursor = db_conn.cursor()
consumer = Consumer(kafka_config)
consumer.subscribe(['customer_orders'])
try:
while True:
msg = consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
if msg.error().code() != KafkaError._PARTITION_EOF:
print(f"Consumer error: {msg.error()}")
continue
payload = json.loads(msg.value().decode('utf-8'))
insert_query = """
INSERT INTO customer_orders (order_id, customer_id, total_amount)
VALUES (%s, %s, %s)
ON CONFLICT (order_id) DO UPDATE
SET total_amount = EXCLUDED.total_amount;
"""
db_cursor.execute(
insert_query,
(
payload['order_id'],
payload['customer_id'],
payload['total_amount']
)
)
db_conn.commit()
# Commit the Kafka offset after the database write succeeds
consumer.commit(msg, asynchronous=False)
except KeyboardInterrupt:
pass
finally:
db_cursor.close()
db_conn.close()
consumer.close()
適したユースケース
カスタムコンシューマー方式は、PostgreSQLへ書き込む前にカスタム変換、データ補完、フィルタリング、条件付きルーティングを必要とするパイプラインに適しています。
最大限の柔軟性が得られる一方、アプリケーションコード、エラーハンドリング、スケーリング、運用監視をチーム側で管理する必要があります。
複雑なイベント処理、ウィンドウ集計、ステートフル処理、ストリーム結合を要する大規模パイプラインでは、単純なコネクタやカスタムコンシューマーでは処理能力が不足する場合があります。
このような場面ではApache FlinkやApache Spark Structured Streamingといったストリーム処理フレームワークが一般的に使用されます。
(Kafkaトピック)—>[Flink/Spark処理エンジン]—>[PostgreSQLデータベース]
JDBCシンクを使用したFlink
Apache Flinkは低レイテンシのストリーム処理向けに設計され、連続するデータストリームに対するステートフル計算をサポートします。FlinkのDataStream APIを使用することで、Kafkaイベントのフィルタリング、変換、データ補完を行った上で、処理結果をPostgreSQLへ書き込めます。
下記JavaサンプルはFlinkのJdbcSinkを用い、処理済みストリームデータをPostgreSQLテーブルへ書き込む例です。
import org.apache.flink.connector.jdbc.JdbcConnectionOptions;
import org.apache.flink.connector.jdbc.JdbcExecutionOptions;
import org.apache.flink.connector.jdbc.JdbcSink;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
public class FlinkPostgresStreamingJob {
public static void main(String[] args) throws Exception {
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Consume and process stream from Kafka
DataStream<OrderEvent> processedStream = env.fromSource(...)
.filter(event -> event.getAmount() > 10.0);
// Write processed data to PostgreSQL
processedStream.addSink(JdbcSink.sink(
"INSERT INTO processed_orders (order_id, amount) VALUES (?, ?) " +
"ON CONFLICT (order_id) DO UPDATE SET amount = EXCLUDED.amount",
(statement, order) -> {
statement.setString(1, order.getOrderId());
statement.setDouble(2, order.getAmount());
},
JdbcExecutionOptions.builder()
.withBatchSize(1000)
.withBatchIntervalMs(200)
.withMaxRetries(5)
.build(),
new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
.withUrl("jdbc:postgresql://postgres-db.internal:5432/ecommerce")
.withDriverName("org.postgresql.Driver")
.withUsername("app_user")
.withPassword("secure_password")
.build()
));
env.execute("Flink-to-PostgreSQL-Sink");
}
}
Spark Structured Streamingが適するケース
Apache Spark Structured Streamingはデフォルトでマイクロバッチ実行モデルを採用しており、Sparkエコシステムを既に利用しているチームに適しています。
Sparkがより適する状況:
適したユースケース
ストリーム処理フレームワークは、高度な変換、ステートフル処理、データ補完、複数のストリームソースの結合を行った上で結果をPostgreSQLへロードする大規模パイプラインに適しています。
本番環境でKafka‑to‑PostgreSQLパイプラインを運用するには、信頼性、パフォーマンス、データ整合性に十分配慮する必要があります。適切な監視、エラーハンドリング、復旧戦略により、データ量や処理の複雑さが増大してもパイプラインを安定に保てます。
コンシューマーラグは、Kafkaパーティション上の最新オフセットと、コンシューマーグループが処理済みのオフセットとの差を示します。ラグが増大し続ける場合、コンシューマーが到着するイベントに追いつけず、データがPostgreSQLに届くまでに遅延が発生します。
コンシューマーラグを低減・管理するための施策:
ストリーミングパイプラインは、不正なメッセージ、スキーマの問題、一時的なデータベース接続異常などの予期せぬ障害に対処しなければなりません。適切なエラーハンドリングがないと、1件の異常レコードによりメッセージ処理が中断する可能性があります。
推奨される施策:
Kafkaのメッセージ保持機能により、アプリケーションのバグ、処理障害、データ同期の問題が発生した際に過去のイベントをリプレイできます。
安全な復旧とリプレイをサポートするため:
ここで紹介した各手法はそれぞれ異なる課題を解決します。スキーマが安定しており、最小限のコードでデータを流したい場合はKafka Connectが適しています。メッセージ処理を細かく制御したい場合にはカスタムコンシューマー、変換、結合、集計処理が必要なパイプラインにはFlinkまたはSparkの複雑さが正当化されます。
いずれの手法も、Kafkaにクリーンでタイムリーなイベントデータが届いていることを前提としています。本番データベースから、遅延、スキーマ不整合、サイレントなデータドリフトなしに、信頼できるデータをKafkaへ取り込むことは別の課題です。Info2softのi2Streamはこの前段階の処理のために開発されています。
i2StreamはKafka関連パイプラインに関連する複数の機能を備えています。
真のボトルネックが、まず本番データベースから信頼できる低レイテンシのデータを取り出すことにあるチームにとって、i2Streamがこのレイヤーを担うことで、下流のKafka Connect、カスタムコンシューマー、Flink/Sparkジョブは信頼できるデータを処理できるようになります。またInfo2softは、ストリーミングから災害復旧へ目的が移った場合に向け、バイトレベルの継続的データ保護を実現するi2CDPも提供しています。
KafkaからPostgreSQLへのデータストリーミングに万能な選択肢は存在しません。スキーマが安定していればKafka Connectが最速で導入でき、処理ロジックに特殊な要件があればカスタムコンシューマーが制御性を提供し、変換処理が必要な場面ではFlinkまたはSparkが力を発揮します。
どの手法を選択しても、パイプラインの信頼性はKafkaに投入される元のデータに依存します。そこでInfo2softのi2Streamのようなツールが活躍し、ソースデータを正確かつ最新の状態に保ち、下流のストリーミングに上流の問題を引き継がせません。
まずはチームの現在のニーズに合致する手法を選択し、パイプラインの複雑さが増大するにつれて選択を見直してください。