使用 Java SDK 为数据表创建全量、增量或全量加增量类型的通道。
前提条件
功能说明
调用 createTunnel 为数据表创建通道。同一数据表可以创建多个通道。通道类型决定消费的数据范围:BaseData 仅消费全量数据,Stream 仅消费增量数据,BaseAndStream 先消费全量数据,再消费增量数据。
创建 Stream 或 BaseAndStream 类型的通道时,如果数据表未开启 Stream,系统会自动开启 Stream,并将增量日志过期时间设置为 7 天。
CreateTunnelResponse createTunnel(CreateTunnelRequest request)
throws TableStoreException, ClientException
以下示例为 example_table 创建全量类型的通道 example_tunnel。
String tableName = "example_table";
String tunnelName = "example_tunnel";
CreateTunnelRequest request =
new CreateTunnelRequest(
tableName, tunnelName, TunnelType.BaseData);
CreateTunnelResponse response = tunnelClient.createTunnel(request);
System.out.println("TunnelId: " + response.getTunnelId());
参数说明
CreateTunnelRequest 包含以下参数:
|
名称 |
类型 |
说明 |
|
tableName(必选) |
|
数据表名称。 |
|
tunnelName(必选) |
|
通道名称。 |
|
tunnelType(必选) |
|
通道类型。 |
|
streamTunnelConfig(可选) |
|
增量数据范围配置,用于 |
|
streamRecordOptions(可选) |
|
增量记录内容配置,用于 |
增量数据范围
streamTunnelConfig 的类型为 StreamTunnelConfig,包含以下参数:
|
名称 |
类型 |
说明 |
|
flag(可选) |
|
未设置 |
|
startOffset(可选) |
|
增量数据的起始时间戳。单位为毫秒,取值范围为 [当前系统时间 - Stream 过期时间 + 5 分钟,当前系统时间)。设置此参数后, |
|
endOffset(可选) |
|
增量数据的结束时间戳,单位为毫秒。同时设置起止时间时,此参数必须大于 |
Stream 过期时间是增量日志的保留时长,最大值为 7 天。为数据表开启 Stream 时可以设置该值,设置后不能修改。
增量记录内容
streamRecordOptions 的类型为 StreamRecordOptions,包含以下参数:
|
名称 |
类型 |
说明 |
|
getVersionGeneratorValue(可选) |
|
增量记录是否包含版本号生成器的值。默认值为 |
|
getSysColumns(可选) |
|
增量记录是否包含系统列。默认值为 |
|
getNewRowInfo(可选) |
|
增量记录是否包含最新行信息。默认值为 |
|
oldColumnsToGet(可选) |
|
原始行中需要返回的属性列。 |
|
newColumnsToGet(可选) |
|
最新行中需要返回的属性列。 |
增量记录列
streamRecordOptions.oldColumnsToGet 和 streamRecordOptions.newColumnsToGet 的类型均为 StreamColumn,包含以下参数:
|
名称 |
类型 |
说明 |
|
columnType(必选) |
|
属性列的选择方式。 |
|
columnNames(可选) |
|
|
返回值
CreateTunnelResponse 包含以下返回字段:
|
字段 |
类型 |
说明 |
|
tunnelId |
|
创建的通道 ID,通过 |
场景示例
指定增量数据范围
以下示例创建增量类型的通道,并指定最近一小时内的增量数据范围。
long endTime = System.currentTimeMillis() - 1_000L;
long startTime = endTime - 60 * 60 * 1000L;
StreamTunnelConfig streamConfig =
new StreamTunnelConfig(startTime, endTime);
CreateTunnelRequest request =
new CreateTunnelRequest(
"example_table",
"example_stream_tunnel",
TunnelType.Stream);
request.setStreamTunnelConfig(streamConfig);
CreateTunnelResponse response = tunnelClient.createTunnel(request);
System.out.println("TunnelId: " + response.getTunnelId());
配置增量记录内容
以下示例创建增量类型的通道,并设置增量记录返回版本号生成器的值、系统列、最新行信息、原始行的 value 属性列和最新行的所有属性列。
StreamColumn oldColumns =
new StreamColumn(StreamColumnType.SPECIFIED_COLUMN);
oldColumns.addColumnName("value");
StreamRecordOptions recordOptions = new StreamRecordOptions();
recordOptions.setGetVersionGeneratorValue(true);
recordOptions.setGetSysColumns(true);
recordOptions.setGetNewRowInfo(true);
recordOptions.setOldColumnsToGet(oldColumns);
recordOptions.setNewColumnsToGet(
new StreamColumn(StreamColumnType.ALL_COLUMNS));
CreateTunnelRequest request =
new CreateTunnelRequest(
"example_table", "example_stream_tunnel", TunnelType.Stream);
request.setStreamRecordOptions(recordOptions);
CreateTunnelResponse response = tunnelClient.createTunnel(request);
System.out.println("TunnelId: " + response.getTunnelId());