You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

JDBC Connector能否将多张表的数据读取到单个Topic中?

如何通过JDBC Connector将多张数据库表数据写入单个Kafka Topic

当然可以实现,核心思路是先用JDBC Source Connector将每张表同步到独立的Kafka Topic,再通过ksqlDB或Kafka Streams将多个Topic的数据合并到单个目标Topic,以下是具体操作步骤:

方法一:使用ksqlDB

  1. 同步单表到独立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。

  2. 在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'
    );
    
  3. 合并流到单个目标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(数据流)

  1. 同步单表到独立Topic
    同方法一,先为每张表配置JDBC Source Connector,将数据同步到各自的Kafka Topic。

  2. 编写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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 19:30:40