Java SDK 可建立投遞任務,將資料表中的全量資料、增量資料或全量與增量資料投遞到同地區的 OSS Bucket。
前提條件
安裝Tablestore Java SDK並初始化用戶端。資料湖投遞功能需要 5.10.3 及以上版本,建議使用最新版本。
開通Object Storage Service,並在 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);