创建同步Paimon

更新时间:
复制 MD 格式

创建同步 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,格式为 oss://<bucket>/<path>

Database

目标 Paimon Database 名称。创建同步任务时,Database 不存在则自动创建。

Table

目标 Paimon 表名。创建同步任务时,表不存在则自动创建;表已存在时校验 Schema 和分区兼容性。

认证方式

支持 DataHub 默认角色、AccessKey 和自定义角色。推荐使用 DataHub 默认角色。

RoleName

使用角色认证时必填

DataHub 默认角色为 AliyunDataHubAccessingOssRole

导入字段

从 Topic Schema 中选择需要写入 Paimon 的字段。字段名和类型应与已存在目标表兼容。

分区模式

支持 SYSTEM_TIMEMETA_TIMEUSER_DEFINE

提交间隔(秒)

Paimon Commit 间隔,默认 300 秒,取值范围为 60~3600 秒。间隔越小,数据可见性越及时,但提交频率越高。

Timestamp Unit

指定 TIMESTAMP 字段的时间单位,支持 MICROSECONDMILLISECONDSECOND,默认为 MICROSECOND

起始位置/起始时间

指定同步任务开始消费 Topic 数据的位置。创建前请确认是否需要同步历史数据。

建表配置

仅当目标表不存在、需要由 DataHub 自动建表时,才需要填写以下配置;目标表已存在时无需填写。

参数

是否必填

说明

分区字段

定义自动创建表的 Paimon 分区字段。

表自定义属性

设置自动创建表的自定义属性,例如 bucket

步骤四:配置分区模式

SYSTEM_TIME

使用 DataHub 记录的系统时间生成 Paimon 时间分区。该模式适合按数据写入 DataHub 的时间组织目标表。

时间分区参数说明如下。

参数

说明

分区配置

配置分区字段及时间格式,例如 ds:%Y%m%dhh:%Hmm:%M。时间生成的分区字段在 Paimon 表中应为 STRING 类型。

分区间隔

时间分桶间隔,默认 15 分钟。

时区

生成时间分区时使用的时区,默认为 Asia/Shanghai

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 应处于 RUNNINGEXECUTING 状态。

  • 同步点位:表示任务已消费到的 DataHub 数据位置。

  • 同步延迟:持续增大时,应检查源 Topic 流量、OSS 访问和任务错误信息。

  • 脏数据量:存在脏数据时,应检查字段类型、NULL 值和分区字段配置。

  • 重启/停止:配置或权限调整后,可按需重启任务。

使用实时计算 Flink 版验证同步结果

创建任务并向 Topic 写入测试数据后,可以在阿里云实时计算 Flink 版(VVP)中查询用户 OSS 上的 Paimon 表,确认数据是否已提交。

  1. 参考管理Paimon Catalog,在实时计算 Flink 版中创建 Paimon Catalog。

  2. 确认 Catalog 配置的 Warehouse 与 DataHub 同步任务中的 Warehouse 完全一致,并具有对应 OSS 路径的访问权限。

  3. 在 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 表,因此故障恢复或重试时可能产生重复记录。对唯一性有要求的业务,应在下游查询或加工链路中根据业务主键去重。