Java SDK 可创建投递任务,将数据表中的全量数据、增量数据或全量与增量数据投递到同地域的 OSS Bucket。
前提条件
安装Tablestore Java SDK并初始化客户端。数据湖投递功能需要 5.10.3 及以上版本,建议使用最新版本。
开通对象存储 OSS,并在 Tablestore 实例所在地域创建 OSS Bucket。具体操作,请参见OSS 快速入门。
创建服务关联角色
AliyunServiceRoleForOTSDataDelivery,并获取角色 ARN。具体操作,请参见创建投递任务。
功能说明
调用 createDeliveryTask 方法创建投递任务。投递任务可以只投递全量数据或增量数据,也可以在完成全量数据投递后持续投递增量数据。
public CreateDeliveryTaskResponse createDeliveryTask(CreateDeliveryTaskRequest request)
throws TableStoreException, ClientException
创建投递任务后,任务需要先完成初始化。可调用 describeDeliveryTask 方法查询投递任务信息。
以下示例创建全量与增量数据投递任务,使用数据写入 Tablestore 的时间按天生成 OSS 分区。示例中的 client 是已初始化的客户端。请将占位符替换为实际值,并确保 pk、event_time 和 active 列的数据类型分别为 String、String 和 Boolean。
String tableName = "<TABLE_NAME>";
String taskName = "<TASK_NAME>";
OSSTaskConfig taskConfig = new OSSTaskConfig();
taskConfig.setOssPrefix("delivery/year=$yyyy/month=$MM/day=$dd");
taskConfig.setOssBucket("<OSS_BUCKET>");
taskConfig.setOssEndpoint("<OSS_ENDPOINT>");
taskConfig.setOssStsRole("<ROLE_ARN>");
taskConfig.addParquetSchema(new ParquetSchema("pk", "pk", DataType.UTF8));
taskConfig.addParquetSchema(
new ParquetSchema("event_time", "event_time", DataType.UTF8));
taskConfig.addParquetSchema(new ParquetSchema("active", "active", DataType.BOOL));
CreateDeliveryTaskRequest request =
new CreateDeliveryTaskRequest(tableName, taskName, taskConfig);
request.setTaskType(DeliveryTaskType.BASE_INC);
client.createDeliveryTask(request);
参数说明
投递请求
request 的类型为 CreateDeliveryTaskRequest,包含以下参数。
|
名称 |
类型 |
说明 |
|
tableName(必选) |
String |
数据表名称。 |
|
taskName(必选) |
String |
投递任务名称。只能包含英文小写字母、数字、短横线( |
|
taskConfig(必选) |
OSSTaskConfig |
OSS 投递配置。 |
|
taskType(必选) |
DeliveryTaskType |
投递任务类型。取值范围: |
OSS 投递配置
request.taskConfig 的类型为 OSSTaskConfig,包含以下参数。
|
名称 |
类型 |
说明 |
|
ossPrefix(必选) |
String |
OSS Bucket 中的目录前缀。支持 |
|
ossBucket(必选) |
String |
OSS Bucket 名称。Bucket 必须与 Tablestore 实例位于同一地域。 |
|
ossEndpoint(必选) |
String |
OSS Bucket 所在地域的服务地址。 |
|
ossStsRole(必选) |
String |
服务关联角色 |
|
parquetSchema(必选) |
List<ParquetSchema> |
待投递字段列表。可选择任意字段并自定义字段投递到 OSS 后的名称和顺序,列表顺序决定字段在 Parquet 文件中的顺序。通过 |
|
eventTimeColumn(可选) |
EventColumn |
事件时间列。设置后, |
|
format(可选) |
OSSFileFormat |
OSS 文件格式。默认值和当前唯一支持的值均为 |
|
timeFormatter(可选) |
TimeFormatter |
预留的分区格式参数。当前版本的 SDK 不会将该参数写入请求,请勿设置。 |
投递字段
request.taskConfig.parquetSchema[] 中每个元素的类型为 ParquetSchema,包含以下参数。
|
名称 |
类型 |
说明 |
|
columnName(必选) |
String |
Tablestore 数据表中的源字段名称。 |
|
ossColumnName(必选) |
String |
字段投递到 OSS 后的名称。 |
|
type(必选) |
DataType |
字段在 Parquet 文件中的目标类型。该类型必须与源字段的数据类型匹配,否则该字段值会作为脏数据丢弃。关于类型映射,请参见数据格式映射。 |
|
encode(可选) |
OSSFileEncoding |
Parquet 编码方式。默认值为 |
|
typeExtend(可选) |
String |
预留的 Parquet 扩展类型参数,当前不支持,请勿设置。 |
事件时间列
request.taskConfig.eventTimeColumn 的类型为 EventColumn,包含以下参数。
|
名称 |
类型 |
说明 |
|
columnName(必选) |
String |
作为事件时间的源字段名称。 |
|
timeFormat(必选) |
EventTimeFormat |
事件时间格式。取值范围: |
场景示例
按事件时间分区
如需根据 event_time 列的时间生成 OSS 分区,请在创建请求前为 taskConfig 设置事件时间列。以下示例要求该列中的时间符合 RFC 3339 格式。
EventColumn eventColumn =
new EventColumn("event_time", EventTimeFormat.RFC3339);
taskConfig.setEventTimeColumn(eventColumn);