Kafka Consumer Integration Task (Beta)
本文介绍如何创建 Kafka Consumer 任务,以持续消费 Kafka topic 中的消息,并将消息内容保存到内部对象存储(租户 Stage)。
与 S3、MySQL 或 PostgreSQL 数据集成任务不同,Kafka Consumer 任务不会直接写入常规目标表。任务创建并启动后,你可以使用 @kafka_consumer/<task_name>/ stage 路径查看已保存的消息对象,并通过 SQL 查询其内容。
如果你需要先创建可复用的 Kafka 连接设置,请参见 Kafka - 凭证(Beta)。
使用场景
- 持续从 Kafka topic 中摄取 JSON 消息
- 先将 Kafka 消息落盘到内部对象存储,再通过下游 SQL 进行查询或处理
- 为实时或准实时数据流水线保留原始 Kafka 消息对象
工作流
- 上游系统将消息写入 Kafka topic。
- Kafka Consumer 任务从指定的 topic 中读取消息。
- 任务将消息批量保存到内部对象存储(租户 Stage)。
- 用户通过
@kafka_consumer/<task_name>/查看生成的对象。 - 用户从 stage 查询消息内容,并根据需要执行下游加载或转换。
前提条件
在创建 Kafka Consumer 任务之前,请确保:
- 已创建 Kafka - Credentials 数据源
- 平台可以通过网络访问 Kafka broker
- Kafka 数据源中的认证方法、TLS 设置和账户信息正确
- Kafka 用户具有读取目标 topic 的权限
- 目标 topic 中的消息与任务中选择的 Data Format 一致
创建 Kafka Consumer 任务
第 1 步:基本信息
进入 Data > Data Integration,然后点击 Create Task。
选择一个 Kafka 数据源,然后配置基本参数:
第 2 步:预览数据
完成基本设置后,点击 Next 进入 Preview Data Info。
系统会尝试从指定的 Kafka topic 中读取示例消息。如果有可用消息,页面会显示 1 到 2 条 JSON 消息,供你验证 topic、数据格式和消息结构。
如果没有可预览的消息,页面会显示 No sample data available。你仍然可以继续创建任务,但建议检查这些 topic 是否已包含消息,以及所选 Start Position 是否能够读取到示例数据。
第 3 步:查看结果
在 Result Viewing 步骤中,选择用于运行 Kafka Consumer 任务的计算集群 (Warehouse)。
任务启动后,会读取 Kafka 消息并将其保存到内部对象存储(租户 Stage)。页面会提供 SQL 示例。你可以使用 LIST @kafka_consumer/<task_name>/ 查看生成的对象,并使用 stage 查询读取消息内容。
-- List stage objects:
LIST @kafka_consumer/<task_name>/;
-- Query object data (replace with the correct PATTERN path):
SELECT $1
FROM @kafka_consumer (
FILE_FORMAT=>'ndjson',
PATTERN=>'<task_name>/year=YYYY/month=MM/day=DD/hour=HH/.*[.]ndjson'
);
点击 Create 创建任务。
任务行为
Kafka Consumer 任务会持续运行。启动后,它会从指定的 topic 中消费消息,并将其批量保存为内部对象存储中的对象文件,直到你手动停止该任务。
查询已保存的消息
Kafka Consumer 任务会将消息对象保存在 @kafka_consumer/<task_name>/ 路径下。任务启动并写入对象后,打开任务详情页并切换到 Data Browsing 页签,即可按 UTC 小时查看对象数量和对象列表。
你也可以先使用 SQL 列出对象,再根据实际路径查询其内容:
LIST @kafka_consumer/<task_name>/;
SELECT $1
FROM @kafka_consumer (
FILE_FORMAT=>'ndjson',
PATTERN=>'<task_name>/year=YYYY/month=MM/day=DD/hour=HH/.*[.]ndjson'
);
如果你需要将消息写入业务表,请基于查询结果继续执行下游转换或加载。
高级配置
运行时大小
Kafka Consumer 任务支持修改运行时大小。在修改 Runtime Size 之前,请先停止任务,然后通过 Edit 菜单打开编辑页面,在 Runtime Size 部分选择合适的运行时大小并保存更改。重启任务后,任务将以新的运行时大小运行。