如何使用Stream从Kafka Topic读取数据及解决KSQL创建报错
问题背景
- 目标:创建KSQL流
MY_STREAM,消费MY_TOPIC的消息并产出到TARGET_TOPIC - 初始执行语句:
CREATE STREAM MY_STREAM WITH (KAFKA_TOPIC='TARGET_TOPIC') AS SELECT * FROM MY_TOPIC;
- 报错信息:
MY_TOPIC doesn't exist - 已知规则:无论Stream采用何种方式创建(是否使用
AS SELECT语法),创建完成后都会和指定Kafka Topic绑定,支持双向的数据推送、拉取。
解决方法
报错和Topic读权限无关,核心原因是KSQL的流创建语法要求,AS SELECT子句的数据源必须是已经在KSQL元数据中注册过的流或表,不能直接读取未注册的原生Kafka Topic,按两步操作即可实现需求:
- 先为源Topic
MY_TOPIC注册对应的源输入流,配置和Topic实际情况匹配的序列化格式、消息Schema、分区数即可,这个操作仅注册元数据映射,不会额外拷贝Topic数据,无额外存储开销:
-- 请将字段、格式参数替换为实际业务配置 CREATE STREAM MY_TOPIC_SOURCE ( -- 示例字段,按实际消息结构填写 -- id VARCHAR, -- op_ts BIGINT, -- biz_content VARCHAR ) WITH ( KAFKA_TOPIC = 'MY_TOPIC', VALUE_FORMAT = 'JSON', -- 支持JSON/AVRO/PROTOBUF/DELTA等,和源Topic消息格式一致 PARTITIONS = 3 -- 和MY_TOPIC实际分区数保持一致即可 );
- 再执行流创建语句,从已注册的源流读取数据,写入绑定
TARGET_TOPIC的MY_STREAM:
CREATE STREAM MY_STREAM WITH ( KAFKA_TOPIC = 'TARGET_TOPIC', VALUE_FORMAT = 'JSON' -- 可按需配置目标Topic的消息序列化格式 ) AS SELECT * FROM MY_TOPIC_SOURCE;
执行完成后,MY_STREAM会持续消费MY_TOPIC的新增消息,写入TARGET_TOPIC,满足需求。
内容的提问来源于stack exchange,提问作者alphanumeric
相关产品推荐
相关产品推荐

