すべてのプロダクト
Search
ドキュメントセンター

Tablestore:トンネルからのデータ消費

最終更新日:Aug 01, 2026

Tablestore SDK for Java は、トンネルからデータを継続的に消費し、各バッチをコールバックで処理します。また、ハートビート、チェックポイント、スレッドプール、および消費の同時実行性を設定できます。

注意事項

  • 増分ログの保持期間は、テーブルのストリームログの有効期限と同じで、最大 7 日間です。 BaseAndStream トンネルの場合、この期間内に全量データ消費が完了しないと、増分データ消費が開始されるときに OTSTunnelExpired エラーが返されます。 このトンネルは、増分データの消費を継続できなくなります。

  • 増分データの消費が遅れて保持期間を過ぎた場合、トンネルは利用可能な最新のデータから再開される場合があります。その結果、一部のデータが消費されない可能性があります。

  • 期限切れのトンネルは無効になる場合があります。トンネルが無効な状態が 30 日以上続いた場合、そのトンネルは削除され、復元できなくなります。

前提条件

Tablestore SDK for Java をインストールして、TunnelClient を初期化します。

機能の説明

TunnelWorker は、トンネル ID によってトンネルに接続し、ハートビートを使用して現在のクライアントに割り当てられたチャネルを取得し、データを継続的にプルして、レコードの各バッチを IChannelProcessor に渡します。複数の TunnelWorker インスタンスが同じトンネルを消費する場合、サーバーはクライアント間でチャネルを分散します。

トンネルからデータを消費するには、次の手順を実行します:

  1. IChannelProcessor を実装します。process メソッドを使用して各バッチを処理し、shutdown メソッドを使用してコールバックで使用されるリソースを解放します。

  2. TunnelWorkerConfig オブジェクトを作成して、コールバックと消費動作を設定します。

  3. トンネル ID、TunnelClient、および TunnelWorkerConfig を使用して TunnelWorker を作成します。

  4. connectAndWorking を呼び出して、消費を開始します。

    void process(ProcessRecordsInput input);
    void shutdown();

次の例では、トンネルからプルされた各レコードを出力し、消費を開始します。

private static class SimpleProcessor implements IChannelProcessor {
    @Override
    public void process(ProcessRecordsInput input) {
        for (StreamRecord record : input.getRecords()) {
            System.out.println(record);
        }
    }

    @Override
    public void shutdown() {
        // コールバックで使用されるリソースを解放します。
    }
}

String tunnelId = "example_tunnel_id";
TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
TunnelWorker worker =
        new TunnelWorker(tunnelId, tunnelClient, config);
worker.connectAndWorking();
重要

connectAndWorking は、バックグラウンドの消費タスクを開始した後に戻ります。アプリケーションプロセスは実行し続けてください。消費を停止するには、worker.shutdown()、config.shutdown()、tunnelClient.shutdown() の順に呼び出します。worker.shutdown() はトンネル接続を閉じ、コールバックの shutdown メソッドを呼び出します。config.shutdown() は、読み取り、処理、およびヘルパースレッドプールをシャットダウンします。TunnelWorker は、ワーカーを停止しようとする JVM シャットダウンフックを登録しますが、アプリケーションは引き続きこれらのリソースを明示的に解放する必要があります。

パラメーター

ワーカー

TunnelWorker コンストラクターには、次のパラメーターが含まれています。

名前

タイプ

説明

tunnelId (必須)

文字列

トンネル ID。トンネルの作成、一覧表示、またはクエリによって取得します。

client (必須)

TunnelClientInterface

初期化された TunnelClient。

workerConfig (必須)

TunnelWorkerConfig

コールバックと消費の動作に関する設定。

消費設定

workerConfig は TunnelWorkerConfig 型で、以下のパラメーターが含まれます。

名前

タイプ

説明

channelProcessor (必須)

IChannelProcessor

データ処理のコールバック。このパラメーターは、3 つのパラメーターを持つ TunnelWorker コンストラクターを使用する場合に必須です。

heartbeatTimeoutInSec (オプション)

long

ハートビートのタイムアウト時間 (秒) です。デフォルト値は 300 で、heartbeatIntervalInSec より大きい必要があります。ハートビートがタイムアウトすると、サーバーはクライアントを利用不可と見なし、クライアントはトンネルに再接続します。

heartbeatIntervalInSec (オプション)

long

ハートビートの間隔 (秒) です。デフォルト値は 30、最小値は 5 です。ハートビートはアクティブなチャネルを取得し、チャネルの状態を更新し、データ処理タスクを初期化します。この間隔は TunnelWorker のウォームアップ時間にも影響します。

checkpointIntervalInMillis (オプション)

long

サーバー上で消費チェックポイントが記録される間隔 (ミリ秒単位) です。デフォルト値は 5000 です。Tunnel Service は、各レコードを少なくとも 1 回配信し、レコードの順序を維持します。再起動されたタスクは最新のチェックポイントから再開するため、一部のデータを複数回処理する可能性があります。間隔を短くすると重複処理は減少しますが、チェックポイントをあまりにも頻繁に記録するとスループットが低下する可能性があります。

clientTag (オプション)

文字列

クライアント ID を生成し、TunnelWorker インスタンスを区別するために使用されるカスタムクライアントタグです。デフォルト値は、Java の os.name システムプロパティです。

readRecordsExecutor (オプション)

ThreadPoolExecutor

データをプルするスレッドプールです。デフォルトのプールには、32 個のコアスレッド、最大 1000 個のスレッド、16 のキュー容量、および 60 秒のキープアライブ時間があります。

processRecordsExecutor (オプション)

ThreadPoolExecutor

データを処理するスレッドプール。そのデフォルト設定はreadRecordsExecutorと同じです。カスタムプールの場合、トンネルのチャネル数に基づいてスレッド数を設定します。

maxChannelParallel (オプション)

int

同時にデータをプルおよび処理するチャネルの最大数です。このパラメーターを使用して、メモリ使用量を制限します。デフォルト値は -1 で、制限がないことを示します。このパラメーターは、Tablestore SDK for Java 5.10.0 以降でサポートされています。

channelHelperExecutor (オプション)

ThreadPoolExecutor

チャネルの初期化、パイプラインのスケジュール、ランタイムエラーの処理を行うヘルパースレッドプール。このパラメーターが設定されていない場合、キャッシュされたスレッドプールが使用されます。

maxRetryIntervalInMillis (オプション)

int

増分データプル中の指数バックオフの最大ベース間隔をミリ秒単位で指定します。デフォルト値は 2000、最小値は 200 です。バッチのレコード数が 500 以下、かつサイズが 900 KB 以下の場合、クライアントはバックオフ間隔を徐々に増やします。実際の間隔は、現在のベース間隔の 75% から 125% の範囲でランダムに選択されます。このパラメーターは、Tablestore SDK for Java 5.4.0 以降でサポートされています。

readMaxTimesPerRound (オプション)

int

パイプラインの1ラウンドにおける ReadRecords 呼び出しの最大数。デフォルト値は 1 です。

readMaxBytesPerRound (オプション)

int

1 回のパイプラインラウンドで取得されるデータの最大量 (バイト単位) で、デフォルト値は 4194304 (4 MiB) です。この値または readMaxTimesPerRound に達すると、ラウンドは停止します。

enableClosingChannelDetect (オプション)

ブール値

CLOSING 状態のチャネルをリアルタイムで検出するかどうかを指定します。CLOSING チャネルは、あるクライアントから別のクライアントへの移行中です。Tablestore SDK for Java 5.13.13 以降がこのパラメーターをサポートしています。バージョン 5.17.0 以降では、デフォルト値は true です。検出を無効にすると、チャネルは多数存在するもののクライアントリソースが不足している場合に、チャネルの移行がブロックされて消費が中断される可能性があります。

同一マシン上で複数の TunnelWorker インスタンスを起動する場合、1 つの TunnelWorkerConfig オブジェクトを再利用して、読み取りおよび処理スレッドプールを共有できます。すべてのワーカーが停止した後、config.shutdown() を 1 回だけ呼び出します。

コールバックデータ

process メソッドは、以下のフィールドを含む ProcessRecordsInput オブジェクトを受け取ります。

フィールド

タイプ

説明

records

List<StreamRecord>

現在のバッチで取得されたレコード。取得するには getRecords() を呼び出します。

nextToken

文字列

次のバッチ用のトークンです。getNextToken() を呼び出して取得します。TunnelWorker はこの値を自動的に使用して、データのプルを続行し、チェックポイントを記録します。

traceId

文字列

現在のプルリクエストのトレース ID。getTraceId() を呼び出して取得します。

channelId

文字列

現在のバッチが属するチャネルの ID です。getChannelId() を呼び出して取得します。チャネル ID からパーティション ID を取得するには、getPartitionId() を呼び出します。

シナリオ例

消費パラメーターのチューニング

消費のスループットやメモリ使用量が要件を満たさない場合は、ハートビートとチェックポイントの間隔、チャネルの同時実行数、ラウンドごとのプル回数とサイズ、および増分プルのバックオフ間隔を調整します。

TunnelWorkerConfig config =
        new TunnelWorkerConfig(new SimpleProcessor());
config.setHeartbeatIntervalInSec(10);
config.setHeartbeatTimeoutInSec(60);
config.setCheckpointIntervalInMillis(10_000);
config.setMaxChannelParallel(16);
config.setReadMaxTimesPerRound(4);
config.setReadMaxBytesPerRound(8 * 1024 * 1024);
config.setMaxRetryIntervalInMillis(3_000);