消息入湖

更新时间:
复制 MD 格式

云消息队列 Kafka 版的消息入湖功能,允许您将 Kafka Topic 中的实时消息自动、持续地写入 OSS Table Bucket(基于 Apache Iceberg 格式),无需额外部署 Flink 或 Spark 作业,即可实现流式数据入湖,构建 Lakehouse(湖仓一体)架构。

功能简介

消息入湖是云消息队列 Kafka 版提供的原生数据入湖能力。开启该功能后,写入 Topic 的消息在 Kafka 集群持久化的基础上,会同时以 Iceberg 表格式写入 OSS Table Bucket,供下游多种计算引擎(MaxCompute、Hologres、Spark、Trino 等)直接查询分析。

您无需自行开发和运维 ETL 流水线,即可实现从实时消息到结构化数据湖的端到端打通。

应用场景

实时数据湖构建

将用户点击流、IoT 传感器数据、交易日志等通过 Kafka 实时写入 OSS 表,数据湖中始终包含最新数据,支持近实时分析。

CDC 数据入湖

通过 Debezium 等工具捕获数据库变更并发送到 Kafka,消息入湖功能以 CDC 模式消费并合并(Upsert)到 Iceberg 表中,实现准实时数据同步。

低成本历史数据归档

Kafka 中短期保留原始事件,OSS 长期存储结构化历史数据。结合 Iceberg 的分区和压缩策略,显著降低存储成本。

Exactly-Once 写入保障

入湖进度维护在 Kafka Leader 元数据中,彻底去除对外挂 KV 等外部系统的依赖,保证了强一致性,同时消除系统间耦合,极大简化了整体链路逻辑,避免数据重复或丢失,满足金融、交易等对数据准确性要求极高的场景。

功能优势

优势

说明

零 ETL 运维

无需单独部署 Flink/Spark 作业,免去数百个流作业的调度、监控、告警和资源管理负担

Schema 统一管理

基于 Kafka Schema Registry 进行统一校验和演进,避免脏数据和格式错误,消除元数据分散问题

自动化数据湖维护

自动处理小文件合并(Compaction)、过期快照清理、分区优化等 Iceberg 表维护任务

成本优势

按实际计算资源(CU)使用量付费,无需预置计算集群,大幅降低数据入湖成本

开放生态

写入 Iceberg 开放格式,支持 MaxCompute、Hologres、Spark、Flink、Trino、Presto 等引擎直接查询

工作原理

消息入湖的数据流转过程如下:

  1. 生产者将消息写入 Kafka Topic。

  2. Kafka 集群对消息进行持久化存储。

  3. 消息入湖模块根据配置的解析模式和写入模式进行数据转换。

  4. 转换后的数据以 Iceberg 表格式写入 OSS Table Bucket。

  5. 下游计算引擎通过 Iceberg Catalog 访问 OSS 中的表数据。

使用限制

限制项

说明

支持版本

仅 Serverless 版本支持消息入湖功能,最低版本要求3.7.0.0

存储要求

需要配置 OSS Table Bucket 作为 Catalog 存储

开服地域

杭州

计费说明

说明

当前入湖能力公测中,正式收费会提前1个月发送公告及通知。

开启消息入湖

使用消息入湖能力,需要首先开启实例入湖能力。在实例开启消息入湖能力后,按照Topic粒度进行数据写入OSS Table Bucket。

前提条件

  • 已创建云消息队列 Kafka 版 Serverless 实例。

操作步骤

  1. 登录云消息队列 Kafka 版控制台

  2. 在左侧导航栏,单击实例列表

  3. 实例列表页面,单击目标实例名称。

  4. 在左侧导航栏,单击消息入湖

  5. 消息入湖页面,按照步骤完成开启消息入湖的相关配置。包含:

    1. Kafka实例版本验证,如低于入湖能力版本,需要完成Kafka实例版本升级;

    2. 跳转OSS Table Bucket页面,创建Table Bucket,详见OSS Tables

    3. 开启Schema Registry;

    4. OSS服务授权;

    5. Catalog关联

    参数

    说明

    Catalog 配置

    添加已创建的 OSS Table Bucket 作为数据湖 Catalog 存储

    服务授权

    授权 Kafka 服务角色对目标 OSS Table Bucket 的读写权限

  6. 完成以上配置后,单击确定,开启消息入湖功能。

执行结果

开启成功后,消息入湖页面显示为已开启状态,您可以开始为 Topic 配置消息入湖规则。

配置 Topic 消息入湖

本节介绍如何为云消息队列 Kafka 版的 Topic 配置消息入湖,将 Topic 中的消息实时写入 OSS Table Bucket(Iceberg 格式)。

前提条件

  • 已开启消息入湖功能。具体操作,请参见开启消息入湖。

  • 目标 Topic 已创建并有消息写入。

操作步骤

  1. 登录云消息队列 Kafka 版控制台。

  2. 在左侧导航栏,单击实例列表

  3. 实例列表页面,单击目标实例名称。

  4. 在左侧导航栏,单击Topic管理

  5. Topic列表页面,单击目标Topic名称。

  6. Topic详情页面,单击消息入湖,按照以下步骤完成设置。

步骤一:配置基础属性

参数

说明

是否必填

Topic 名称

入湖的目标 Topic

命名空间(Namespace)

Iceberg 表的命名空间,用于在 Catalog 中进行逻辑分组

表名称

Iceberg 表名,默认与 Topic 名称一致,不可修改

是(自动填充)

表分区类型

Iceberg 表的分区策略,用于优化查询性能。表分区类型可参见属性说明/表分区类型。

步骤二:配置写入模式

参数

说明

可选值

解析模式

消息写入 Iceberg 的方式

原始归档、结构化归档

写入模式

数据写入 Iceberg 表的模式

Append:追加;CDC:变更数据捕获。

valueConvertType(Value 转换类型)

消息 Value 的反序列化方式

raw:原始字节;string:字符串;by_schema_id:按消息中嵌入的 Schema ID 解析

keyConvertType(Key 转换类型)

消息 Key 的反序列化方式

raw:原始字节;string:字符串;by_schema_id:按消息中嵌入的 Schema ID 解析

TransformType(转换类型)

消息体结构转换方式

none: 不变换;flatten: 展平嵌套结构;flatten_debezium: 处理 Debezium CDC 事件(需要 schema-based 转换)

写入模式为CDC 时的附加配置:

选择 CDC 写入模式时,需额外配置:

参数

说明

可选值

CDCType(CDC 类型)

CDC 消息的处理方式

通用 CDC 字段、Upsert 模式

IdColumns(主键列)

Iceberg 表的主键列,用于 CDC 操作的行级别定位

自定义输入

  • 通用 CDC 字段:需要填写CDCField。

  • Upsert 模式:根据主键列进行 Upsert(存在则更新,不存在则插入)操作。

参数

说明

可选值

CDCField

根据消息中的用来标记CDC操作类型的字段判断执行 Insert、Update 或 Delete 操作。

自定义输入

步骤三:配置同步策略

参数

说明

默认值

提交间隔

数据写入 Iceberg 表的提交周期,单位为毫秒(ms),取值范围 60,000ms ~ 900,000ms

60,000ms

错误容忍

遇到异常数据时的处理策略

invalid_data

错误容忍策略说明:

取值

说明

none

不容忍任何错误,遇到异常数据立即停止入湖任务

invalid_data

仅容忍无效数据错误,跳过格式异常的消息并继续处理

all

容忍所有错误,跳过所有异常并继续处理

完成配置

配置完成后,单击确定,创建 Topic 消息入湖配置。

执行结果

配置创建成功后,系统将自动开始消费 Topic 消息并写入 OSS Table Bucket 。您可以在消息入湖页面查看入湖状态和运行指标,同时可以从Table Bucket中读取数据。

消息入湖监控指标

本节介绍云消息队列 Kafka 版消息入湖功能支持的监控指标,帮助您了解入湖任务的运行状态和性能表现。

Topic 级别指标

指标名称

指标类型

说明

Topic 转 Table Time Lag 延迟

Gauge

表示该 Topic 目标表中已写入数据的 Kafka timestamp 与当前时间的差距。计算方式为:当前时间戳 - Topic 写入进度时间戳;其中每个分区的写入进度时间戳取该分区已写入目标表数据里的最大 Kafka record timestamp,Topic 级别再取所有有效分区中的最小值。因此 Topic 级别等价于所有有效分区该时间差的最大值,单位 ms。该指标反映表内数据时间相对当前时间的落后程度,不等同于消费 offset lag。如果该值持续增大,可能是某些分区入湖处理进度落后,也可能是源端持续写入 timestamp 本身较旧的数据,需要结合 offset lag、写入吞吐和消息 timestamp 分布一起判断。

Topic 转 Table Offset Lag 延迟

Counter

该 Topic 当前还有多少数据没有完成入湖。计算方式:Topic 最新数据位置 - 当前已经成功写入目标表的数据位置。Topic 级别会把所有分区的滞后量加和。反映入湖堆积量。它比 Time Lag 更偏底层,关注的是"还剩多少条数据没追上"。如果该值持续增大,说明入湖速度低于生产速度,需要关注处理能力、写入能力或资源瓶颈。

入湖消费 Topic 的吞吐

Gauge

入湖任务从 Topic 读取数据的速度,单位: Bytes/s。衡量入湖任务从源 Topic 拉取数据的能力。如果这个值长期低于业务写入 Topic 的速度,后续就会产生堆积,Offset Lag 会逐步升高。

入湖写入 OSS 的吞吐

Gauge

该 Topic 的数据被转换并写入 OSS 的速度,单位:Bytes/s。衡量数据真正落到存储上的能力。这个指标过低或波动很大时,可能说明 OSS 写入、网络带宽、文件提交批次、或者本地写入策略存在瓶颈。

入湖处理 Kafka Record 的速度(Message/s)

Gauge

该 Topic 的数据成功推进入湖进度的速度,也就是每秒有多少条记录完成处理并被纳入目标表进度,单位:Message/s。衡量整体入湖处理能力,包括读取、解析、转换、写入和提交进度等环节。这个值越高,说明单位时间内完成入湖的数据越多;如果该值偏低,而资源使用率较高,通常需要进一步排查处理逻辑或写入链路瓶颈。

属性说明

表分区类型说明

支持 7 种分区类型,以下逐一介绍每种类型的语义、适用场景和配置示例。

Identity(原始值分区)

按字段原始值直接分区。配置中直接写字段名,不带任何 transform 函数。

语法: [column_name]

配置示例:

region

分区效果:

region 值

分区路径

cn-hangzhou

region=cn-hangzhou/

us-east-1

region=us-east-1/

eu-west-1

region=eu-west-1/

适用场景: 地域、环境、业务线等低基数离散字段。

注意事项: 字段基数直接等于分区数。高基数字段(如用户ID)会导致分区爆炸,应使用 bucket 代替。

year(按年分区)

提取 timestamp/date 类型字段的年份部分。

语法: [year(column)]

配置示例:

year(create_time)

分区效果:

create_time 值

分区路径

2026-04-27 10:30:00

create_time_year=2026/

2025-12-01 08:00:00

create_time_year=2025/

2024-06-15 12:00:00

create_time_year=2024/

适用场景: 数据跨度大、按年归档或查询的长周期数据(如历史订单、审计日志)。

month(按月分区)

提取 timestamp/date 类型字段的年月部分。

语法: [month(column)]

配置示例:

month(event_time)

分区效果:

event_time 值

分区路径

2026-04-27 10:30:00

event_time_month=2026-04/

2026-03-15 08:00:00

event_time_month=2026-03/

适用场景: 最常用的时间分区方式。适合日志、事件流、指标数据等按月维度查询的场景。每年只产生 12 个分区,元数据增长可控。

day(按天分区)

提取 timestamp/date 类型字段的年月日部分。

语法: [day(column)]

配置示例:

day(log_time)

分区效果:

log_time 值

分区路径

2026-04-27 10:30:00

log_time_day=2026-04-27/

2026-04-28 08:00:00

log_time_day=2026-04-28/

适用场景: 大数据量日志分析、运营报表等按天维度查询的场景。每年产生 365 个分区,需配合 snapshot expiration 清理历史分区。

hour(按小时分区)

提取 timestamp 类型字段的年月日时部分。

语法: [hour(column)]

配置示例:

hour(ingest_time)

分区效果:

ingest_time 值

分区路径

2026-04-27 10:30:00

ingest_time_hour=2026-04-27-10/

2026-04-27 11:15:00

ingest_time_hour=2026-04-27-11/

适用场景: 超高吞吐场景,查询需要精确到小时。每年产生 8760 个分区,元数据增长快,注意:

  • 每小时单分区数据量需足够大(至少数十 MB),否则小文件问题严重

  • 建议配合 RewriteDataFiles 定期合并

bucket(哈希桶分区)

对字段值做 Murmur3 哈希,分配到 N 个固定桶中。适用于高基数字段的均匀分布。

语法: [bucket(column, N)],N 为桶数(正整数)

配置示例:

# 按 user_id 哈希到 16 个桶
bucket(user_id, 16)

分区效果:

user_id 值

分区路径

12345

user_id_bucket=7/(hash(12345) % 16)

67890

user_id_bucket=2/(hash(67890) % 16)

11111

user_id_bucket=14/(hash(11111) % 16)

适用场景: 用户ID、订单ID、设备ID 等高基数字段,需要分散数据但无法直接用 identity 的场景。

桶数 N 选择建议:

数据量级(单次查询范围)

建议 N

每天 < 1GB

4~8

每天 1~10GB

8~16

每天 10~100GB

16~64

每天 > 100GB

64~128

注意: N 一旦设定不可直接修改(修改配置会触发 Partition Evolution,旧数据仍在旧桶数下,新数据用新桶数)。

truncate(截断分区)

对字段值按宽度 W 截断,将连续值映射到离散的分区。

语法: [truncate(column, W)],W 为截断宽度(正整数)

对不同数据类型的截断行为:

字段类型

截断规则

示例(W=100)

int / long

value - (value % W)

1234 → 1200, 567 → 500, 50 → 0

string

取前 W 个字符

W=3 时 "hangzhou" → "han"

decimal

类似整数,按精度截断

-

配置示例(整数截断):

truncate(price, 100)

分区效果:

price 值

分区路径

1234

price_trunc=1200/

567

price_trunc=500/

50

price_trunc=0/

99

price_trunc=0/

配置示例(字符串截断):

# 按 city_name 前 3 个字符分区
truncate(city_name, 3)

分区效果:

city_name 值

分区路径

hangzhou

city_name_trunc=han/

shanghai

city_name_trunc=sha/

shenzhen

city_name_trunc=she/

适用场景: 数值范围查询(价格区间、年龄段、分数段)、字符串前缀聚合(城市名、产品编码前缀)。

组合分区

多个分区字段可以组合使用,逗号分隔,Iceberg 按声明顺序构建多级分区:

# 时间 + 哈希桶(最常见的组合)
[day(event_time), bucket(user_id, 16)]

# 地域 + 按月
[region, month(create_time)]

# 三级分区
[year(ts), bucket(tenant_id, 8), region]

组合分区的目录结构示例(时间 + 哈希桶):

data/${hash}/
  event_time_day=2026-04-27/
    user_id_bucket=0/
      00001.parquet
    user_id_bucket=1/
      00002.parquet
  event_time_day=2026-04-28/
    user_id_bucket=0/
      00003.parquet

组合原则:

  • 查询最常用的过滤条件放前面(利于分区裁剪)

  • 避免组合后分区总数过大(如 day × bucket(64) = 每天 64 个分区,一年 23360 个)

  • 低基数字段用 identity,高基数字段用 bucket

解析模式说明

  • 原始归档:将消息以原始字节或字符串形式直接存入 Iceberg 表,适用于日志归档、原始数据备份等场景。写入模式固定为 Append,TransformType 固定为 none。

  • 结构化归档:基于 Schema 对消息进行解析和结构化处理后写入 Iceberg 表,支持 Append 和 CDC 两种写入模式,适用于需要结构化查询和分析的场景。

写入模式说明

  • Append:追加写入模式,所有消息直接追加至 Iceberg 表,适用于日志、事件流等只增数据场景。

  • CDC:变更数据捕获模式,支持根据消息中的操作类型字段进行 Insert / Update / Delete 操作,适用于数据库变更同步入湖场景。