如何使用Confluent Schema Registry在KSQL中创建可关联的users主键表
用Confluent Schema Registry在KSQL中创建可关联的users表
前提确认
- 确保KSQL已连接到Confluent Schema Registry,且
users主题的Schema已在Registry中注册(从你提供的Key结构来看,这是Debezium捕获MySQL变更的典型格式,一般Key和Value都会有对应的Schema) - 确认
users主题的Value部分Schema包含user_id字段(或确认Key中的id字段与user_id对应)
步骤1:查看主题结构(可选但推荐)
先在KSQL中打印主题的第一条消息,确认完整的字段结构:
PRINT 'users' FROM BEGINNING LIMIT 1;
这会输出消息的Key和Value Schema,帮你确认主键字段的名称、类型等细节。
步骤2:创建KSQL表
根据你提供的Key Schema(id为string类型),分两种场景创建表:
场景1:主键来自消息Key
如果user_id就是消息Key中的id字段,直接关联Key创建表:
CREATE TABLE users ( user_id STRING PRIMARY KEY, -- 按实际Value Schema添加其他字段,比如 name STRING, email STRING 等 ) WITH ( KAFKA_TOPIC = 'users', VALUE_FORMAT = 'AVRO', -- 若使用JSON Schema则改为'JSON_SCHEMA' KEY_FORMAT = 'AVRO', PARTITIONS = 6, -- 替换为你的主题实际分区数 REPLICAS = 3 -- 替换为你的主题实际副本数 );
场景2:主键来自消息Value
如果user_id存在于消息Value中,指定KEY_FIELD来关联:
CREATE TABLE users ( user_id STRING PRIMARY KEY, -- 按实际Value Schema添加其他字段 ) WITH ( KAFKA_TOPIC = 'users', VALUE_FORMAT = 'AVRO', -- 若使用JSON Schema则改为'JSON_SCHEMA' KEY_FORMAT = 'AVRO', PARTITIONS = 6, REPLICAS = 3, KEY_FIELD = 'user_id' -- 指定Value中的字段作为主键 );
处理Debezium的变更消息(补充)
如果users主题是Debezium捕获的MySQL变更,消息Value会包含before和after字段,需要提取最新状态的字段:
- 先创建流:
CREATE STREAM users_stream WITH ( KAFKA_TOPIC = 'users', VALUE_FORMAT = 'AVRO', KEY_FORMAT = 'AVRO' );
- 基于流创建表,提取
after中的最新数据:
CREATE TABLE users AS SELECT after->id AS user_id, after->name AS name, after->email AS email FROM users_stream EMIT CHANGES;
步骤3:验证表可用性
创建完成后,查询表确认数据:
SELECT * FROM users LIMIT 10;
此时该表即可通过user_id与其他KSQL表/流进行JOIN关联操作。
内容的提问来源于stack exchange,提问作者ashanr
相关产品推荐
相关产品推荐

