创建同步 Paimon 任务
DataHub 支持将 Tuple Topic 中的数据实时同步到 Apache Paimon 表,Paimon Warehouse 存储在用户自己的阿里云 OSS Bucket 中。本文介绍如何在 DataHub 控制台创建、配置和验证 Paimon 同步任务。
邀测说明
Paimon 同步功能目前处于邀测阶段。如需使用,请通过阿里云工单提交申请,并提供阿里云账号 UID。
当前邀测版本暂不支持带主键的 Paimon 表,仅支持无主键的 Append 表。
功能限制
源 Topic 必须为 Tuple 类型,暂不支持 Blob Topic。
目标 Warehouse 必须位于用户自己的 OSS Bucket 中,格式为
oss://<bucket>/<path>。当前邀测版本暂不支持带主键的 Paimon 表。
当前邀测版本暂不支持同步 DataHub JSON 类型字段。
同步任务采用至少一次(At-least-once)语义。目标表为 Append 表时,异常重试等场景可能产生重复数据,下游应根据业务需要进行去重。
DataHub Topic 与 OSS Bucket 必须位于同一 Region。
准备工作
1. 申请邀测资格
请提交阿里云工单申请邀测,并提供阿里云账号 UID。
2. 准备 DataHub Topic
创建 Tuple Topic,并确认需要同步的字段已定义在 Topic Schema 中。Paimon 同步暂不支持 Blob Topic。
DataHub 与 Paimon 的字段类型对应关系如下。
DataHub 类型 | Paimon 类型 | 是否支持 |
BOOLEAN | BOOLEAN | 支持 |
TINYINT | TINYINT | 支持 |
SMALLINT | SMALLINT | 支持 |
INTEGER | INT | 支持 |
BIGINT | BIGINT | 支持 |
FLOAT | FLOAT | 支持 |
DOUBLE | DOUBLE | 支持 |
DECIMAL | DECIMAL | 支持 |
STRING | STRING | 支持 |
TIMESTAMP | TIMESTAMP | 支持 |
ARRAY | ARRAY | 支持 |
MAP | MAP | 支持 |
STRUCT | ROW/STRUCT | 支持 |
JSON | JSON | 当前邀测版本暂不支持 |
3. 准备目标 OSS 和 Paimon 表
准备用户自己的 OSS Bucket 和 Warehouse 路径,例如 oss://example-bucket/paimon/warehouse。目标 Database 和 Paimon 表可以预先创建;如果不存在,创建同步任务时会自动创建。
创建任务时,控制台会自动检测 AliyunDataHubAccessingOssRole 是否存在;如果不存在,可以按照页面提示一键创建并授权访问目标 OSS。
Warehouse 是用户 OSS 中的 Paimon 根目录,不是 OSS Endpoint,也不是 DataHub 提供的存储空间。
创建同步任务
步骤一:进入同步任务入口
登录 DataHub 控制台,进入目标 Project 和 Topic,在 Topic 详情页单击右上角的 同步。
步骤二:选择 Paimon
在 新建 Connector 面板中选择 Paimon。
步骤三:配置目标 Paimon 表
填写 Warehouse、Database、Table、认证方式和导入字段等参数。
主要参数说明如下。
参数 | 是否必填 | 说明 |
Warehouse | 是 | 用户自己的 OSS Paimon Warehouse,格式为 |
Database | 是 | 目标 Paimon Database 名称。创建同步任务时,Database 不存在则自动创建。 |
Table | 是 | 目标 Paimon 表名。创建同步任务时,表不存在则自动创建;表已存在时校验 Schema 和分区兼容性。 |
认证方式 | 是 | 支持 DataHub 默认角色、AccessKey 和自定义角色。推荐使用 DataHub 默认角色。 |
RoleName | 使用角色认证时必填 | DataHub 默认角色为 |
导入字段 | 是 | 从 Topic Schema 中选择需要写入 Paimon 的字段。字段名和类型应与已存在目标表兼容。 |
分区模式 | 是 | 支持 |
提交间隔(秒) | 是 | Paimon Commit 间隔,默认 300 秒,取值范围为 60~3600 秒。间隔越小,数据可见性越及时,但提交频率越高。 |
Timestamp Unit | 否 | 指定 TIMESTAMP 字段的时间单位,支持 |
起始位置/起始时间 | 是 | 指定同步任务开始消费 Topic 数据的位置。创建前请确认是否需要同步历史数据。 |
建表配置
仅当目标表不存在、需要由 DataHub 自动建表时,才需要填写以下配置;目标表已存在时无需填写。
参数 | 是否必填 | 说明 |
分区字段 | 否 | 定义自动创建表的 Paimon 分区字段。 |
表自定义属性 | 否 | 设置自动创建表的自定义属性,例如 |
步骤四:配置分区模式
SYSTEM_TIME
使用 DataHub 记录的系统时间生成 Paimon 时间分区。该模式适合按数据写入 DataHub 的时间组织目标表。
时间分区参数说明如下。
参数 | 说明 |
分区配置 | 配置分区字段及时间格式,例如 |
分区间隔 | 时间分桶间隔,默认 15 分钟。 |
时区 | 生成时间分区时使用的时区,默认为 |
META_TIME
优先使用记录携带的 Meta Time 生成时间分区;记录未携带 Meta Time 时使用 DataHub 系统时间。适用于需要保留上游元数据时间的场景。
META_TIME 的分区配置、分区间隔和时区配置方式与 SYSTEM_TIME 相同。
USER_DEFINE
使用 Topic 中与目标 Paimon 表分区字段同名的业务字段作为分区值。仅当目标表不存在、需要由 DataHub 自动建表时,才需要在 建表配置 > 分区字段 中选择对应字段;如果目标表已存在,则无需选择,系统会按照已有表的分区定义从 Topic 记录中获取分区值。
使用 USER_DEFINE 时请注意:
自动建表时,选择的业务分区字段必须存在于 Topic Schema 中。
目标表已存在时,其业务分区字段必须存在于 Topic Schema 中,但无需在 建表配置 > 分区字段 或普通导入字段中重复选择。
写入记录的分区字段不能为 NULL。
如果目标表已存在,字段名称、字段类型和分区字段顺序应与目标表兼容。
USER_DEFINE 不使用时间格式化分区配置。
步骤五:创建任务
确认配置无误后,单击 创建。控制台会检查 OSS 访问权限、Paimon Database、目标表、Schema 和分区配置。
目标 Database 不存在时,控制台自动创建 Database。
目标表不存在时,控制台自动创建 Paimon 表。
目标表已存在时,控制台复用该表并执行兼容性校验。
如果目标表带主键,当前邀测版本会创建失败。
查看和管理任务
进入 Topic 详情页的 同步任务 页签,选择 Paimon 任务查看运行详情。
重点关注以下信息:
任务状态:任务和各 Shard 对应 Task 应处于
RUNNING或EXECUTING状态。同步点位:表示任务已消费到的 DataHub 数据位置。
同步延迟:持续增大时,应检查源 Topic 流量、OSS 访问和任务错误信息。
脏数据量:存在脏数据时,应检查字段类型、NULL 值和分区字段配置。
重启/停止:配置或权限调整后,可按需重启任务。
使用实时计算 Flink 版验证同步结果
创建任务并向 Topic 写入测试数据后,可以在阿里云实时计算 Flink 版(VVP)中查询用户 OSS 上的 Paimon 表,确认数据是否已提交。
参考管理Paimon Catalog,在实时计算 Flink 版中创建 Paimon Catalog。
确认 Catalog 配置的 Warehouse 与 DataHub 同步任务中的 Warehouse 完全一致,并具有对应 OSS 路径的访问权限。
在 SQL 开发页面执行以下查询:
USE CATALOG `<catalog-name>`;
USE `<database>`;
SET 'execution.runtime-mode' = 'batch';
SHOW TABLES;
SELECT *
FROM `<table>`
LIMIT 10;也可以使用完整表名直接查询:
SELECT *
FROM `<catalog-name>`.`<database>`.`<table>`
LIMIT 10;如果使用时间分区,可以增加分区条件,例如:
SELECT *
FROM `<table>`
WHERE ds = '20260810'
AND hh = '15'
LIMIT 10;常见问题
目标 Database 和表必须提前创建吗?
不必须。目标 Paimon Database 和表可以提前创建;如果不存在,创建同步任务时会自动创建。
是否支持带主键的 Paimon 表?
当前邀测版本暂不支持。请使用无主键的 Append 表;后续支持情况以产品公告和正式文档为准。
数据是否可能重复?
同步任务采用至少一次语义,目标又是无主键 Append 表,因此故障恢复或重试时可能产生重复记录。对唯一性有要求的业务,应在下游查询或加工链路中根据业务主键去重。