创建通道

更新时间:
复制 MD 格式

使用 Java SDK 为数据表创建全量、增量或全量加增量类型的通道。

前提条件

功能说明

调用 createTunnel 为数据表创建通道。同一数据表可以创建多个通道。通道类型决定消费的数据范围:BaseData 仅消费全量数据,Stream 仅消费增量数据,BaseAndStream 先消费全量数据,再消费增量数据。

说明

创建 StreamBaseAndStream 类型的通道时,如果数据表未开启 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(必选)

String

数据表名称。

tunnelName(必选)

String

通道名称。

tunnelType(必选)

TunnelType

通道类型。BaseData 表示仅消费全量数据;Stream 表示仅消费增量数据;BaseAndStream 表示先消费全量数据,再消费增量数据。

streamTunnelConfig(可选)

StreamTunnelConfig

增量数据范围配置,用于 StreamBaseAndStream 类型的通道。如果设置此参数,时间范围必须有效。

streamRecordOptions(可选)

StreamRecordOptions

增量记录内容配置,用于 StreamBaseAndStream 类型的通道。此参数需要使用 5.17.11 及以上版本的 Java SDK。

增量数据范围

streamTunnelConfig 的类型为 StreamTunnelConfig,包含以下参数:

名称

类型

说明

flag(可选)

StartOffsetFlag

未设置 startOffset 时的增量数据起点。LATEST 表示从通道创建时间开始,EARLIEST 表示从当前可读取的最早增量日志开始。默认值为 LATEST

startOffset(可选)

long

增量数据的起始时间戳。单位为毫秒,取值范围为 [当前系统时间 - Stream 过期时间 + 5 分钟,当前系统时间)。设置此参数后,flag 不生效。

endOffset(可选)

long

增量数据的结束时间戳,单位为毫秒。同时设置起止时间时,此参数必须大于 startOffset。不设置时持续消费增量数据。

说明

Stream 过期时间是增量日志的保留时长,最大值为 7 天。为数据表开启 Stream 时可以设置该值,设置后不能修改。

增量记录内容

streamRecordOptions 的类型为 StreamRecordOptions,包含以下参数:

名称

类型

说明

getVersionGeneratorValue(可选)

boolean

增量记录是否包含版本号生成器的值。默认值为 false

getSysColumns(可选)

boolean

增量记录是否包含系统列。默认值为 false

getNewRowInfo(可选)

boolean

增量记录是否包含最新行信息。默认值为 false

oldColumnsToGet(可选)

StreamColumn

原始行中需要返回的属性列。

newColumnsToGet(可选)

StreamColumn

最新行中需要返回的属性列。

增量记录列

streamRecordOptions.oldColumnsToGetstreamRecordOptions.newColumnsToGet 的类型均为 StreamColumn,包含以下参数:

名称

类型

说明

columnType(必选)

StreamColumnType

属性列的选择方式。SPECIFIED_COLUMN 表示指定属性列,INPUT_COLUMNS 表示本次写入或更新的属性列,ALL_COLUMNS 表示所有属性列。

columnNames(可选)

List<String>

columnTypeSPECIFIED_COLUMN 时需要返回的属性列名称列表。

返回值

CreateTunnelResponse 包含以下返回字段:

字段

类型

说明

tunnelId

String

创建的通道 ID,通过 getTunnelId() 获取。消费通道数据时需要使用该 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());