使用 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());