JDBC Connector能否将多张表的数据读取到单个Topic中?
如何通过JDBC Connector将多张数据库表数据写入单个Kafka Topic
当然可以实现,核心思路是先用JDBC Source Connector将每张表同步到独立的Kafka Topic,再通过ksqlDB或Kafka Streams将多个Topic的数据合并到单个目标Topic,以下是具体操作步骤:
方法一:使用ksqlDB
同步单表到独立Topic
为每张数据库表配置单独的JDBC Source Connector,示例配置片段:name=jdbc-source-table1 connector.class=io.confluent.connect.jdbc.JdbcSourceConnector connection.url=jdbc:mysql://db-host:3306/db-name connection.user=db-user connection.password=db-pass table.whitelist=table1 topics=table1_topic mode=incrementing incrementing.column.name=id value.converter=io.confluent.connect.avro.AvroConverter value.converter.schema.registry.url=http://schema-registry:8081复制上述配置,修改
name、table.whitelist、topics参数,为每张表创建对应的Connector。在ksqlDB中创建流
连接ksqlDB CLI后,为每个Topic创建对应的流(假设数据格式为Avro):CREATE STREAM table1_stream (id INT, col1 STRING) WITH ( KAFKA_TOPIC='table1_topic', VALUE_FORMAT='AVRO', KEY_FORMAT='KAFKA' ); CREATE STREAM table2_stream (id INT, col2 INT) WITH ( KAFKA_TOPIC='table2_topic', VALUE_FORMAT='AVRO', KEY_FORMAT='KAFKA' );合并流到单个目标Topic
根据业务需求选择合并方式:- 无关联合并(所有记录聚合):用
UNION ALL将两个流的所有记录写入同一个Topic,需统一字段结构(用CAST(null AS ...)填充缺失字段):CREATE STREAM combined_stream WITH (KAFKA_TOPIC='combined_topic') AS SELECT 'table1' AS source_table, id, col1 AS string_col, CAST(null AS INT) AS int_col FROM table1_stream UNION ALL SELECT 'table2' AS source_table, id, CAST(null AS STRING) AS string_col, col2 AS int_col FROM table2_stream; - 关联合并(按主键匹配):用
JOIN关联两个流的记录(需设置时间窗口避免无限等待):CREATE STREAM joined_stream WITH (KAFKA_TOPIC='combined_topic') AS SELECT t1.id, t1.col1, t2.col2 FROM table1_stream t1 INNER JOIN table2_stream t2 WITHIN 1 HOURS ON t1.id = t2.id;
- 无关联合并(所有记录聚合):用
方法二:使用Kafka Streams(数据流)
同步单表到独立Topic
同方法一,先为每张表配置JDBC Source Connector,将数据同步到各自的Kafka Topic。编写Kafka Streams应用
以Java为例,编写应用读取多个输入Topic,转换为统一格式后合并写入单个Topic:import org.apache.kafka.streams.*; import org.apache.kafka.streams.kstream.*; import io.confluent.kafka.streams.serdes.avro.SpecificAvroSerde; import java.util.*; public class TableCombinerApp { public static void main(String[] args) { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "table-combiner-app"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker:9092"); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.Integer().getClass()); // 配置Avro Serde Map<String, String> avroConfig = new HashMap<>(); avroConfig.put("schema.registry.url", "http://schema-registry:8081"); SpecificAvroSerde<Table1Record> table1Serde = new SpecificAvroSerde<>(); table1Serde.configure(avroConfig, false); SpecificAvroSerde<Table2Record> table2Serde = new SpecificAvroSerde<>(); table2Serde.configure(avroConfig, false); SpecificAvroSerde<CombinedRecord> combinedSerde = new SpecificAvroSerde<>(); combinedSerde.configure(avroConfig, false); StreamsBuilder builder = new StreamsBuilder(); // 读取两个输入流 KStream<Integer, Table1Record> table1Stream = builder.stream( "table1_topic", Consumed.with(Serdes.Integer(), table1Serde) ); KStream<Integer, Table2Record> table2Stream = builder.stream( "table2_topic", Consumed.with(Serdes.Integer(), table2Serde) ); // 转换为统一格式并合并 KStream<Integer, CombinedRecord> combinedStream = table1Stream.mapValues( record -> new CombinedRecord("table1", record.getId(), record.getCol1(), null) ).merge(table2Stream.mapValues( record -> new CombinedRecord("table2", record.getId(), null, record.getCol2()) )); // 写入目标Topic combinedStream.to("combined_topic", Produced.with(Serdes.Integer(), combinedSerde)); KafkaStreams streams = new KafkaStreams(builder.build(), props); Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); streams.start(); } }说明:需提前定义
Table1Record、Table2Record、CombinedRecord的Avro Schema,并通过Schema Registry管理;若需关联数据,可替换merge为join方法并设置时间窗口。
内容的提问来源于stack exchange,提问作者Pooja Srivastava
相关产品推荐
相关产品推荐

