Loading...

We've detected that your browser language is Chinese. Would you like to visit our Chinese website? [ Dismiss ]
By: Emma

現代のイベント駆動型アーキテクチャでは、KafkaからPostgreSQLへのデータストリーミングは一般的な要件です。Kafkaはスケーラブルなリアルタイムイベント処理を実現し、PostgreSQLはアプリケーション、分析、データ同期向けに信頼性の高いリレーショナルストレージを提供します。

本ガイドでは、KafkaのデータがPostgreSQLへ流れる仕組み、一般的な連携手法、信頼性の高いリアルタイムデータパイプラインを構築するための主要なベストプラクティスを解説します。

「KafkaからPostgreSQLへのストリーミング」とは実際何を意味するのか

リアルタイムデータパイプラインにおいて、Apache Kafkaはイベントストリーミングプラットフォームとして動作し、PostgreSQLは処理済みデータを保存・クエリするリレーショナルデータベースとして機能します。

KafkaからPostgreSQLへのデータストリーミングとは、Kafkaトピックからイベントを継続的に消費し、低レイテンシでPostgreSQLテーブルへ書き込むことを指します。

この連携は主に以下の用途で利用されます。

  • リアルタイム分析:トランザクションイベントをPostgreSQLに投入し、リアルタイムダッシュボードやレポーティングを支えます。
  • マテリアライズドビュー:新しいイベントが到着するたびに、事前計算済みの指標や集計データを最新の状態に保ちます。
  • イベント駆動型アプリケーション:上流サービスが生成したイベントに基づき、アプリケーションデータベースを更新します。
  • データ統合パイプライン:複数システムからのイベントストリームを、クエリ可能なリレーショナルデータに変換します。

Kafkaエコシステムにおいて、ソースコネクタはデータをKafkaへ送り込み、シンクコネクタはKafkaからPostgreSQLなど外部システムへデータを転送します。

Kafkaからのデータストリーミングは、データベース同期や高可用性シナリオでよく使われるCDCやデータベースレプリケーションとは異なります。

連携パターン Kafka→PG ストリーミング CDC DBレプリケーション
目的

Kafkaイベントデータを

PostgreSQLテーブルに書き込む

データベースの変更をKafka

または他システムへストリームする

環境間でデータベースのコピー

を維持する

データ単位

イベントメッセージ(JSON、Avro、

Protobuf)

INSERT/UPDATE/DELETEイベント トランザクションまたはログレコード
主なユースケース

リアルタイム分析、

データ統合

リアルタイム同期、イベント駆動アーキテクチャ 災害復旧、高可用性、読み取りスケーリング

KafkaトピックからPostgreSQLテーブルへデータが流れる仕組み

信頼性の高いKafka‑to‑PostgreSQLパイプラインは、Kafkaトピックからリレーショナルテーブルへデータを移送する複数のコンポーネントで構成されます。このデータフローを理解することで、Kafka Connectやカスタムコンシューマーといった各連携手法の裏側の動作を把握できます。

how data moves from kafka topics to postgresql tables

Kafkaプロデューサーとトピックの役割

Kafkaプロデューサーは、Kafkaトピックへレコードを発行するアプリケーションです。これらのレコードはキー‑バリュー形式のメッセージとして保存され、通常JSON、Apache Avro、Protocol Buffersなどの形式でシリアライズされます。

Kafkaトピックはパーティションに分割され、スケーラビリティと並列処理を実現します。プロデューサーがメッセージを送信すると、Kafkaはメッセージキーに基づいてパーティションに割り当てます。同一キーのメッセージは同一パーティションに書き込まれ、処理順序が保たれます。

KafkaコンシューマーによるPostgreSQLへのデータ書き込み

受信側では、コンシューマーまたは連携コネクタがKafkaトピックからメッセージを読み取り、PostgreSQLテーブルへ書き込みます。代表的な処理フローは以下の通りです。

  1. 読み取り:Kafkaトピックから新しいレコードを取得
  2. デシリアライズ:メッセージを構造化データ形式に変換
  3. マッピング:メッセージのフィールドをPostgreSQLテーブルのカラムに対応付け
  4. 書き込み:INSERTまたはUPSERTなどのSQL操作でデータを登録

PostgreSQLへの書き込み処理はKafkaからの読み取りよりもレイテンシが高くなる傾向があるため、スループットを維持しオーバーヘッドを抑えるには、複数レコードをバッチ処理してデータベースへ書き込むことが重要です。

KafkaからPostgreSQLへデータをストリーミングする3つの手法

KafkaからPostgreSQLへデータをストリーミングする一般的な手法は3種類あります。適切な選択は、データ量、変換要件、レイテンシの期待値、運用の複雑さなどの要因に依存します。

手法1:JDBCシンクコネクタを使用したKafka Connect

Kafka ConnectはApache Kafkaエコシステムのオープンソースフレームワークで、Kafkaと外部システム間のデータ移送を実現します。

JDBCシンクコネクタは設定ベースのコネクタで、カスタムコンシューマーコードを作成せずにKafkaトピックからレコードを消費し、PostgreSQLテーブルへ書き込みます。

method 1 kafka connect with jdbc sink connector

JDBCシンクコネクタの動作

JDBCシンクコネクタはKafka Connectクラスター内で動作し、設定されたKafkaトピックからレコードを継続的に消費します。コネクタ設定に基づいてKafkaレコードをデータベース操作に変換し、JDBCを介してPostgreSQLへデータを書き込みます。

コネクタはKafkaコンバーターを通じJSON、Avroなど対応フォーマットでシリアライズされたデータを処理可能です。Avroのようなスキーマベースのフォーマットでは、スキーマレジストリと連携しスキーマ互換性を管理するのが一般的です。

コネクタ設定例

下記はKafkaトピックの注文データをPostgreSQLへ書き込む代表的なJDBCシンクコネクタの設定例です。

json
{
  "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"
  }
}
  • connection.url:PostgreSQLデータベースへ接続するJDBC接続文字列を指定
  • insert.mode:レコードの書き込み方式を定義。upsertを指定すると、主キーが一致する既存行を更新し、重複キーエラーを回避可能
  • pk.mode:コネクタが主キーを特定する方法。record_keyではKafkaメッセージキーをデータベースの主キーとして使用
  • auto.createおよびauto.evolve:対応環境において、レコードのスキーマをもとにテーブル作成やカラム追加を自動実行
注記:auto.evolveは開発・テストを簡略化できますが、本番環境では慎重に使用してください。上流のデータ構造が適切なレビューなしに変更された場合、予期せぬデータベース変更が発生する可能性があります。

適したユースケース

JDBCシンクコネクタによるKafka Connectは、スキーマが比較的安定した単純なデータ取り込みパイプラインに適しています。カスタムコンシューマーの作成・保守をせず、Kafka‑to‑PostgreSQLのデータ移送を安定的に実現したいチームに向いています。

手法2:カスタムコンシューマー(Python/Go)

データ変換、フィルタリング、ルーティング、アプリケーション固有のロジックを完全に制御したい場合、カスタムコンシューマーは柔軟な選択肢となります。既製のコネクタを使用せず、PythonやGoなどの言語で軽量アプリケーションを作成し、Kafkaメッセージを消費してPostgreSQLへ書き込みます。

(Kafkaトピック)—>[カスタムコンシューマーアプリケーション(Python/Go)]—>[PostgreSQLデータベース]

Pythonによるデータ消費と書き込み

カスタムコンシューマーはKafkaクラスターに直接接続し、新規レコードを取得し、アプリケーション層でメッセージを処理した後、結果をPostgreSQLへ書き込みます。

下記の例はconfluent‑kafka PythonクライアントでJSONメッセージを消費し、psycopg2を使用してupsert処理でPostgreSQLへレコードを書き込むサンプルです。

python
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へ書き込む前にカスタム変換、データ補完、フィルタリング、条件付きルーティングを必要とするパイプラインに適しています。

最大限の柔軟性が得られる一方、アプリケーションコード、エラーハンドリング、スケーリング、運用監視をチーム側で管理する必要があります。

手法3:ストリーム処理フレームワーク(Flink/Spark)

複雑なイベント処理、ウィンドウ集計、ステートフル処理、ストリーム結合を要する大規模パイプラインでは、単純なコネクタやカスタムコンシューマーでは処理能力が不足する場合があります。

このような場面ではApache FlinkやApache Spark Structured Streamingといったストリーム処理フレームワークが一般的に使用されます。

(Kafkaトピック)—>[Flink/Spark処理エンジン]—>[PostgreSQLデータベース]

JDBCシンクを使用したFlink

Apache Flinkは低レイテンシのストリーム処理向けに設計され、連続するデータストリームに対するステートフル計算をサポートします。FlinkのDataStream APIを使用することで、Kafkaイベントのフィルタリング、変換、データ補完を行った上で、処理結果をPostgreSQLへ書き込めます。

下記JavaサンプルはFlinkのJdbcSinkを用い、処理済みストリームデータをPostgreSQLテーブルへ書き込む例です。

java
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がより適する状況:

  • チームが既にApache SparkまたはDatabricksのワークロードを運用している
  • Python(PySpark)とDataFrame APIを使用してパイプラインを構築したい
  • 連続的な低レイテンシ処理ではなく、マイクロバッチのレイテンシで許容できるユースケースである

適したユースケース

ストリーム処理フレームワークは、高度な変換、ステートフル処理、データ補完、複数のストリームソースの結合を行った上で結果をPostgreSQLへロードする大規模パイプラインに適しています。

信頼性の高いKafka‑to‑PostgreSQLストリーミングのベストプラクティス

本番環境でKafka‑to‑PostgreSQLパイプラインを運用するには、信頼性、パフォーマンス、データ整合性に十分配慮する必要があります。適切な監視、エラーハンドリング、復旧戦略により、データ量や処理の複雑さが増大してもパイプラインを安定に保てます。

Kafkaコンシューマーラグを監視する

コンシューマーラグは、Kafkaパーティション上の最新オフセットと、コンシューマーグループが処理済みのオフセットとの差を示します。ラグが増大し続ける場合、コンシューマーが到着するイベントに追いつけず、データがPostgreSQLに届くまでに遅延が発生します。

コンシューマーラグを低減・管理するための施策:

  • コンシューマーグループのスケーリング:ワークロード要件とKafkaパーティション数に応じてコンシューマーインスタンスを追加。1つのパーティションは同一コンシューマーグループ内の1つのコンシューマーのみがアクティブに処理可能。
  • バッチ処理の最適化:コンシューマーの取得設定とデータベース書き込みバッチを調整し、スループットとリソース使用量のバランスを取る。小さすぎるバッチはデータベースのオーバーヘッドを増加させ、過大なバッチはメモリ使用量と処理遅延を引き起こす。
  • パイプラインパフォーマンスの監視:Prometheus、Grafana、Burrowなどの監視ツールを使用してコンシューマーラグを追跡し、異常遅延に対するアラートを設定する。

エラーハンドリングとリトライ機構を設計する

ストリーミングパイプラインは、不正なメッセージ、スキーマの問題、一時的なデータベース接続異常などの予期せぬ障害に対処しなければなりません。適切なエラーハンドリングがないと、1件の異常レコードによりメッセージ処理が中断する可能性があります。

推奨される施策:

  • デッドレターキュー(DLQ)の実装:処理失敗したメッセージを別のKafkaトピックへ振り分け、後から調査・再処理を実施。不正レコードがメインパイプラインをブロックすることを防ぐ。
  • リトライ戦略の活用:ネットワーク切断やデータベースの一時的な利用不可といった一時的な障害に対し、指数バックオフ付きのリトライ機構を適用。
  • リカバリ可能エラーと不可能エラーの分離:一時的な障害と、不正なデータフォーマットやスキーマ不整合などの恒久的な障害を区別して処理。失敗レコードをログに記録・追跡することでトラブルシューティングを簡略化。

データ復旧とリプレイを考慮した設計

Kafkaのメッセージ保持機能により、アプリケーションのバグ、処理障害、データ同期の問題が発生した際に過去のイベントをリプレイできます。

安全な復旧とリプレイをサポートするため:

  • 冪等な書き込み処理を設計:同一イベントを複数回安全に処理できるデータベース操作を使用。PostgreSQLのUPSERT(ON CONFLICT DO UPDATE)はリプレイ時の重複レコード防止に有効。
  • オフセットを慎重に管理:必要に応じてコンシューマーグループのオフセットをリセットし、特定のポイントからデータをリプレイ。メッセージをリプレイする前に、既存のPostgreSQLデータへの影響を検証し、書き込み処理が重複イベントに対応できることを確認。

自身のパイプラインに適した手法の選択

ここで紹介した各手法はそれぞれ異なる課題を解決します。スキーマが安定しており、最小限のコードでデータを流したい場合はKafka Connectが適しています。メッセージ処理を細かく制御したい場合にはカスタムコンシューマー、変換、結合、集計処理が必要なパイプラインにはFlinkまたはSparkの複雑さが正当化されます。

いずれの手法も、Kafkaにクリーンでタイムリーなイベントデータが届いていることを前提としています。本番データベースから、遅延、スキーマ不整合、サイレントなデータドリフトなしに、信頼できるデータをKafkaへ取り込むことは別の課題です。Info2softのi2Streamはこの前段階の処理のために開発されています。

i2StreamはKafka関連パイプラインに関連する複数の機能を備えています。

  • ログベースのリアルタイムキャプチャ:テーブルのポーリングではなくデータベースログを直接読み取り、高同時実行環境でもミリ秒レベルのレイテンシを実現。下流システムへ供給するデータを本番環境に追従させ続けます。
  • DDL/DML同期を統合:データ変更と同時にスキーマ変更もレプリケートされるため、ソース側のテーブル変更により、古い構造を前提としたコネクタやコンシューマーが暗黙的に破損することを防ぎます。
  • 組み込みのデータ完全性チェック:MD5チェックサム比較、可視化されたドリフト分析、ワンクリック修復により、カスタム検証スクリプトなしでソースとターゲット間の不整合を自動的に検知します。
  • エージェントレスデプロイ:本番データベースへソフトウェアをインストールする必要がなく、レプリケーション処理がソースシステムのパフォーマンスに一切影響を与えません。
  • 幅広いデータベース・プラットフォーム対応:i2Streamは40種類以上のデータベースおよびビッグデータ環境に対応しており、パイプラインが複数のソースシステムからデータを取得する必要が生じた場合に有用です。

真のボトルネックが、まず本番データベースから信頼できる低レイテンシのデータを取り出すことにあるチームにとって、i2Streamがこのレイヤーを担うことで、下流のKafka Connect、カスタムコンシューマー、Flink/Sparkジョブは信頼できるデータを処理できるようになります。またInfo2softは、ストリーミングから災害復旧へ目的が移った場合に向け、バイトレベルの継続的データ保護を実現するi2CDPも提供しています。

60日間無料トライアル

まとめ

KafkaからPostgreSQLへのデータストリーミングに万能な選択肢は存在しません。スキーマが安定していればKafka Connectが最速で導入でき、処理ロジックに特殊な要件があればカスタムコンシューマーが制御性を提供し、変換処理が必要な場面ではFlinkまたはSparkが力を発揮します。

どの手法を選択しても、パイプラインの信頼性はKafkaに投入される元のデータに依存します。そこでInfo2softのi2Streamのようなツールが活躍し、ソースデータを正確かつ最新の状態に保ち、下流のストリーミングに上流の問題を引き継がせません。

まずはチームの現在のニーズに合致する手法を選択し、パイプラインの複雑さが増大するにつれて選択を見直してください。

概要は準備中です

関連記事

目次:
最新情報を購読
最新のインサイト、ニュース、限定コンテンツをお届けします。いつでも配信解除が可能です。
購読する
ビジネスデータのセキュリティ強化を始めませんか?
60日間の無料トライアルまたはデモで、Info2softが企業データをどのように保護するかをご確認ください。
フォームにご記入の上、送信してください。担当者より追ってご連絡いたします。
このフォームを送信することにより、 プライバシー通知を読み、同意したことを確認します。
{{ isSubmitting ? '送信中...' : '送信する' }}