将数据从 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 数据库复制
目的

将 Kafka 事件数据写入

PostgreSQL 表。

将数据库变更流式传输到 Kafka

或其他系统。

跨环境维护数据库副本。

数据单元

事件消息(JSON、Avro、

Protobuf)。

插入/更新/删除事件。 事务或日志记录。
主要使用场景

实时分析和

数据集成。

实时同步、事件驱动架构。 灾难恢复、高可用、读扩展。

数据如何从 Kafka 主题移动到 PostgreSQL 表

一个可靠的 Kafka 到 PostgreSQL 管道涉及多个组件,这些组件将数据从 Kafka 主题移动到关系型表中。理解这种数据流有助于解释不同的集成方法(如 Kafka Connect 和自定义消费者)在后台是如何工作的。

数据如何从 kafka 主题移动到 postgresql 表

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 有三种常见方法。正确的选择取决于数据量、转换需求、延迟期望和运维复杂性等因素。

方法 1:使用 JDBC Sink Connector 的 Kafka Connect

Kafka Connect 是 Apache Kafka 生态系统中的一个开源框架,支持 Kafka 与外部系统之间的数据移动。

JDBC Sink Connector 是一个基于配置的连接器,它从 Kafka 主题消费记录并将其写入 PostgreSQL 表,无需编写自定义消费者代码。

方法 1 使用 jdbc sink connector 的 kafka connect

JDBC Sink Connector 如何工作

JDBC Sink Connector 在 Kafka Connect 集群中运行,并持续从配置的 Kafka 主题消费记录。它根据连接器配置将 Kafka 记录转换为数据库操作,然后使用 JDBC 将数据写入 PostgreSQL。

该连接器可以通过 Kafka 转换器处理以 JSON、Avro 或其他支持格式序列化的数据。对于基于模式(如 Avro)的格式,通常与 Schema Registry 集成以管理模式兼容性。

连接器配置示例

以下示例展示了一个典型的 JDBC Sink Connector 配置,用于将订单数据从 Kafka 主题写入 PostgreSQL:

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.createauto.evolve:允许连接器在受支持时根据记录模式自动创建表或添加列。
注意:尽管 auto.evolve 可以简化开发和测试,但在生产环境中应谨慎使用。如果上游数据结构在未经适当审查的情况下发生变化,自动模式变更可能会引入意外的数据库修改。

最适合

使用 JDBC Sink Connector 的 Kafka Connect 适用于模式相对稳定的直接数据摄取管道。它非常适合那些需要可靠的 Kafka 到 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 数据库]

使用 Flink 和 JDBC Sink

Apache Flink 专为低延迟流处理而设计,支持对连续数据流进行有状态计算。借助 Flink 的 DataStream API,团队可以在将处理结果写入 PostgreSQL 之前对 Kafka 事件进行过滤、转换和丰富。

以下 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 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 到 PostgreSQL 流式传输最佳实践

在生产环境中运行 Kafka 到 PostgreSQL 管道需要仔细关注可靠性、性能和数据一致性。适当的监控、错误处理和恢复策略有助于在数据量和处理复杂度增加时保持管道稳定。

监控 Kafka 消费者滞后

消费者滞后衡量 Kafka 分区中最新可用偏移量与消费者组已处理偏移量之间的差异。当滞后持续增长时,消费者无法跟上 incoming 事件的速度,导致数据到达 PostgreSQL 之前出现延迟。

减少和管理消费者滞后的方法:

  • 扩展消费者组: 根据工作负载要求和 Kafka 分区数量添加消费者实例。每个分区在消费者组内只能被一个消费者主动处理。
  • 优化批处理: 调整消费者获取设置和数据库写入批次,以平衡吞吐量和资源使用。小批量可能增加数据库开销,而过大的批量可能增加内存使用和处理延迟。
  • 监控管道性能: 使用 Prometheus、Grafana 或 Burrow 等监控工具跟踪消费者滞后,并为异常延迟设置告警。

设计错误处理和重试机制

流式管道必须处理意外故障,包括格式错误的消息、模式问题和临时数据库连接问题。如果没有适当的错误处理,单条有问题的记录就可能中断消息处理。

推荐做法包括:

  • 实现死信队列(DLQ): 将失败的消息路由到单独的 Kafka 主题,以便后续调查和重新处理。这可以防止无效记录阻塞主管道。
  • 使用重试策略: 对于临时性故障(如网络中断或数据库暂时不可用),应用带指数退避的重试机制。
  • 区分可恢复和不可恢复错误: 对临时问题和永久性故障(如无效数据格式或模式不匹配)进行不同处理。记录和跟踪失败记录有助于简化故障排除。

规划数据恢复和重放

Kafka 的消息保留能力允许团队在从应用程序错误、处理失败或数据同步问题中恢复时重放历史事件。

为了支持安全的恢复和重放:

  • 设计幂等写入: 使用可以安全地多次处理同一事件的数据库操作。PostgreSQL 的 UPSERT 操作(如 ON CONFLICT DO UPDATE)可以帮助防止重放场景中的重复记录。
  • 谨慎管理偏移量: 在需要时重置消费者组偏移量以从特定点重放数据。在重放消息之前,验证对现有 PostgreSQL 数据的影响,并确保写入操作被设计为处理重复事件。

为您的管道选择正确的方法

这里介绍的每种方法都解决不同的问题。当您的模式稳定且只需要数据流动而无需编写大量代码时,Kafka Connect 效果很好。当您需要对消息处理进行细粒度控制时,自定义消费者更有意义。当管道需要真正的转换、连接或聚合时,Flink 或 Spark 才值得其复杂性。

这三种方法都假设 Kafka 已经拥有干净、及时的事件数据。以一种可靠的方式从实时生产数据库中将数据获取到 Kafka,而不会出现延迟、模式不匹配或静默漂移,这是另一个问题。英方软件的 i2Stream 正是为这个更早的阶段而构建的。

i2Stream 具有与 Kafka 相关管道相关的几个功能:

  • 基于日志的实时捕获: i2Stream 直接从数据库日志读取,而不是轮询表,在高并发环境中也能实现毫秒级延迟。这使得输入下游系统的数据保持最新,而不是落后于生产环境。
  • 集成 DDL/DML 同步: 模式变更与数据变更一起复制,因此源端的表结构更改不会静默地破坏等待旧结构的连接器或消费者。
  • 内置数据完整性检查: MD5 校验和比较、可视化漂移分析和一键修复可自动捕获源和目标之间的不一致,无需自定义验证脚本。
  • 无代理部署: 生产数据库上无需安装软件,因此复制对源系统的性能产生零影响。
  • 广泛的数据库和平台支持: i2Stream 覆盖 40+ 种数据库和大数据环境,当管道最终需要从多个源系统拉取数据时非常有用。

对于团队来说,如果真正的瓶颈是首先从生产数据库中获取可靠、低延迟的数据,那么 i2Stream 处理这一层,使得下游的 Kafka Connect、自定义消费者或 Flink/Spark 作业有可靠的数据可用。此外,英方软件还提供 i2CDP,用于当目标从流式传输转向灾难恢复时的持续字节级数据保护。

结论

将数据从 Kafka 流式传输到 PostgreSQL 并非一刀切的决策。当您的模式稳定时,Kafka Connect 能最快地将您带到目的地;当逻辑变得具体时,自定义消费者给您控制权;当转换是工作的一部分时,Flink 或 Spark 占有一席之地。

无论您选择哪种方法,管道的可靠性都取决于最初输入 Kafka 的数据。这就是像 英方软件的 i2Stream 这样的工具发挥作用的地方——保持源数据的准确性和时效性,使下游流式传输不会继承上游的问题。

从符合您团队当前需求的方法开始,并随着管道复杂度的增长重新审视这个选择。

博客分类底部

准备好构建企业数据韧性了吗?

立即开启 60 天免费试用,或预约产品演示,了解英方软件如何为您的核心业务提供「零中断、零丢失」的数据保护。

请先完成图形验证

验  证  码:

英方官网验证码
第三方二维码 第三方二维码
英方公告铃铛图标
英方公告铃铛图标

公告

英方侧边栏向右箭头
英方高亮提示圆点
英方软件公告
各位求职者、合作伙伴:
近期有第三方冒用英方名义发布虚假招聘、不实业务信息。我司正规招聘全程零收费,非官网渠道信息均不作数。
信息核验热线:400-0078-655
遇诈骗请保留证据,及时联系我们并报警
英方软件
2026 年 6 月 23 日
英方邮件咨询图标
英方邮件咨询图标

邮件

英方销售支持图标
英方销售支持图标

销售

英方侧边栏向右箭头
联系销售:400-0078-655 转 1