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

如何使用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字段,需要提取最新状态的字段:

  1. 先创建流:
CREATE STREAM users_stream WITH (
  KAFKA_TOPIC = 'users',
  VALUE_FORMAT = 'AVRO',
  KEY_FORMAT = 'AVRO'
);
  1. 基于流创建表,提取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 22:55:31