VMware无法从ISO引导:完整故障排查指南
2026-08-31
2026-08-31
2026-08-31
2026-08-31
将数据从 Kafka 流式传输到 PostgreSQL 是现代事件驱动架构中的常见需求。Kafka 支持可扩展的实时事件处理,而 PostgreSQL 为应用程序、分析和数据同步提供可靠的关系型存储。
本指南涵盖了 Kafka 数据如何流入 PostgreSQL、常见的集成方法以及构建可靠实时数据管道的关键实践。
在实时数据管道中,Apache Kafka 充当事件流平台,而 PostgreSQL 则作为存储和查询已处理数据的关系型数据库。
将数据从 Kafka 流式传输到 PostgreSQL 意味着持续从 Kafka 主题消费事件,并以低延迟将其写入 PostgreSQL 表。
这种集成通常用于:
在 Kafka 生态系统中,源连接器将数据移入 Kafka,而目标连接器将数据从 Kafka 移出到外部系统,如 PostgreSQL。
将数据从 Kafka 流式传输不同于 CDC 和数据库复制,后者通常用于数据库同步和高可用性场景。
| 集成模式 | Kafka→PG 流式 | CDC | 数据库复制 |
|---|---|---|---|
| 目的 |
将 Kafka 事件数据写入 PostgreSQL 表。 |
将数据库变更流式传输到 Kafka 或其他系统。 |
跨环境维护数据库副本。 |
| 数据单元 |
事件消息(JSON、Avro、 Protobuf)。 |
插入/更新/删除事件。 | 事务或日志记录。 |
| 主要使用场景 |
实时分析和 数据集成。 |
实时同步、事件驱动架构。 | 灾难恢复、高可用、读扩展。 |
一个可靠的 Kafka 到 PostgreSQL 管道涉及多个组件,这些组件将数据从 Kafka 主题移动到关系型表中。理解这种数据流有助于解释不同的集成方法(如 Kafka Connect 和自定义消费者)在后台是如何工作的。

Kafka 生产者是将记录发布到 Kafka 主题的应用程序。这些记录以键值消息的形式存储,通常以 JSON、Apache Avro 或 Protocol Buffers 等格式进行序列化。
Kafka 主题被划分为多个分区,以支持可扩展性和并行处理。当生产者发送消息时,Kafka 根据消息键将其分配到某个分区。具有相同键的消息被写入同一分区,这有助于维护它们的处理顺序。
在接收端,消费者或集成连接器从 Kafka 主题读取消息,并将其写入 PostgreSQL 表。典型的工作流程包括:
由于 PostgreSQL 写入通常比 Kafka 读取具有更高的延迟,因此在写入数据库之前对多条记录进行批处理对于维持吞吐量和减少开销非常重要。
将数据从 Kafka 流式传输到 PostgreSQL 有三种常见方法。正确的选择取决于数据量、转换需求、延迟期望和运维复杂性等因素。
Kafka Connect 是 Apache Kafka 生态系统中的一个开源框架,支持 Kafka 与外部系统之间的数据移动。
JDBC Sink Connector 是一个基于配置的连接器,它从 Kafka 主题消费记录并将其写入 PostgreSQL 表,无需编写自定义消费者代码。

JDBC Sink Connector 如何工作
JDBC Sink Connector 在 Kafka Connect 集群中运行,并持续从配置的 Kafka 主题消费记录。它根据连接器配置将 Kafka 记录转换为数据库操作,然后使用 JDBC 将数据写入 PostgreSQL。
该连接器可以通过 Kafka 转换器处理以 JSON、Avro 或其他支持格式序列化的数据。对于基于模式(如 Avro)的格式,通常与 Schema Registry 集成以管理模式兼容性。
连接器配置示例
以下示例展示了一个典型的 JDBC Sink Connector 配置,用于将订单数据从 Kafka 主题写入 PostgreSQL:
{
"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 Sink Connector 的 Kafka Connect 适用于模式相对稳定的直接数据摄取管道。它非常适合那些需要可靠的 Kafka 到 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 数据库]
使用 Flink 和 JDBC Sink
Apache Flink 专为低延迟流处理而设计,支持对连续数据流进行有状态计算。借助 Flink 的 DataStream API,团队可以在将处理结果写入 PostgreSQL 之前对 Kafka 事件进行过滤、转换和丰富。
以下 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 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 到 PostgreSQL 管道需要仔细关注可靠性、性能和数据一致性。适当的监控、错误处理和恢复策略有助于在数据量和处理复杂度增加时保持管道稳定。
消费者滞后衡量 Kafka 分区中最新可用偏移量与消费者组已处理偏移量之间的差异。当滞后持续增长时,消费者无法跟上 incoming 事件的速度,导致数据到达 PostgreSQL 之前出现延迟。
减少和管理消费者滞后的方法:
流式管道必须处理意外故障,包括格式错误的消息、模式问题和临时数据库连接问题。如果没有适当的错误处理,单条有问题的记录就可能中断消息处理。
推荐做法包括:
Kafka 的消息保留能力允许团队在从应用程序错误、处理失败或数据同步问题中恢复时重放历史事件。
为了支持安全的恢复和重放:
这里介绍的每种方法都解决不同的问题。当您的模式稳定且只需要数据流动而无需编写大量代码时,Kafka Connect 效果很好。当您需要对消息处理进行细粒度控制时,自定义消费者更有意义。当管道需要真正的转换、连接或聚合时,Flink 或 Spark 才值得其复杂性。
这三种方法都假设 Kafka 已经拥有干净、及时的事件数据。以一种可靠的方式从实时生产数据库中将数据获取到 Kafka,而不会出现延迟、模式不匹配或静默漂移,这是另一个问题。英方软件的 i2Stream 正是为这个更早的阶段而构建的。
i2Stream 具有与 Kafka 相关管道相关的几个功能:
对于团队来说,如果真正的瓶颈是首先从生产数据库中获取可靠、低延迟的数据,那么 i2Stream 处理这一层,使得下游的 Kafka Connect、自定义消费者或 Flink/Spark 作业有可靠的数据可用。此外,英方软件还提供 i2CDP,用于当目标从流式传输转向灾难恢复时的持续字节级数据保护。
将数据从 Kafka 流式传输到 PostgreSQL 并非一刀切的决策。当您的模式稳定时,Kafka Connect 能最快地将您带到目的地;当逻辑变得具体时,自定义消费者给您控制权;当转换是工作的一部分时,Flink 或 Spark 占有一席之地。
无论您选择哪种方法,管道的可靠性都取决于最初输入 Kafka 的数据。这就是像 英方软件的 i2Stream 这样的工具发挥作用的地方——保持源数据的准确性和时效性,使下游流式传输不会继承上游的问题。
从符合您团队当前需求的方法开始,并随着管道复杂度的增长重新审视这个选择。
公告
邮件
销售