如何使用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
相关产品推荐
相关产品推荐

