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

如何使用KSQL将源Kafka主题拆分为4个专属子主题

使用KSQL拆分嵌套Kafka主题为独立主题的实现方法

前提说明

假设源Kafka主题名为source_topic,消息格式为JSON,包含subscriber、patient、case、service四个嵌套节点。以下是具体拆分步骤:


1. 注册源主题为KSQL流

先将源主题导入KSQL,创建对应的流结构,可选择手动定义Schema(更精准)或自动推断:

方式1:手动定义Schema

CREATE STREAM source_stream (
    subscriber STRUCT<
        mem_ID STRING,
        mem_PIN STRING,
        case_NUMBER STRING,
        mem_FIRST_NAME STRING,
        mem_MIDDLE_NAME STRING,
        mem_LAST_NAME STRING,
        mem_ADD_1 STRING,
        mem_ADD_2 STRING,
        mem_CITY STRING
    >,
    patient STRUCT<
        pat_ID STRING,
        pat_SEX STRING,
        pat_DOB STRING,
        case_NUMBER STRING,
        pat_FIRST_NAME STRING,
        pat_MIDDLE_NAME STRING,
        pat_LAST_NAME STRING,
        pat_PLANE_TYPE STRING,
        pat_PLAN_NAME STRING
    >,
    `case` STRUCT<
        CASE_NUMBER STRING,
        CASE_TYPE STRING,
        CASE_CODE STRING,
        CASE_START_DATE STRING,
        CASE_END_DATE STRING,
        CASE_AUTH_TYPE STRING,
        CASE_STATUS STRING
    >,
    service STRUCT<
        svc_ID STRING,
        case_NUMBER STRING,
        svc_TYPE STRING,
        svc_CODE STRING,
        svc_FAC_ID STRING,
        svc_FAC_NAME STRING,
        svc_PHY_ID STRING,
        svc_PHY_NAME STRING
    >
) WITH (
    KAFKA_TOPIC='source_topic',
    VALUE_FORMAT='JSON',
    KEY_FORMAT='KAFKA'
);

方式2:自动推断Schema

如果源主题已有消息,KSQL可自动识别结构:

SET 'auto.offset.reset' = 'earliest';
CREATE STREAM source_stream WITH (KAFKA_TOPIC='source_topic', VALUE_FORMAT='JSON');

2. 拆分并输出到独立主题

分别创建4个流,提取每个嵌套节点的内容并输出到对应Kafka主题:

拆分到Subscriber主题

CREATE STREAM subscriber_stream WITH (
    KAFKA_TOPIC='Subscriber',
    VALUE_FORMAT='JSON'
) AS SELECT subscriber.* FROM source_stream EMIT CHANGES;

拆分到Patient主题

CREATE STREAM patient_stream WITH (
    KAFKA_TOPIC='Patient',
    VALUE_FORMAT='JSON'
) AS SELECT patient.* FROM source_stream EMIT CHANGES;

拆分到Case主题

注意case是KSQL关键字,需用反引号包裹:

CREATE STREAM case_stream WITH (
    KAFKA_TOPIC='Case',
    VALUE_FORMAT='JSON'
) AS SELECT `case`.* FROM source_stream EMIT CHANGES;

拆分到Service主题

CREATE STREAM service_stream WITH (
    KAFKA_TOPIC='Service',
    VALUE_FORMAT='JSON'
) AS SELECT service.* FROM source_stream EMIT CHANGES;

3. 验证结果

可通过以下方式验证输出内容:

  • KSQL查询示例:
    SELECT * FROM subscriber_stream EMIT CHANGES LIMIT 1;
    
  • Kafka命令行消费示例:
    kafka-console-consumer.sh --bootstrap-server <kafka-broker:port> --topic Subscriber --from-beginning --property print.key=false --property value.deserializer=org.apache.kafka.common.serialization.StringDeserializer
    

内容的提问来源于stack exchange,提问作者Renu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 22:40:49