Tablestore SDK for Java は、トンネルからデータを継続的に消費し、各バッチをコールバックで処理します。また、ハートビート、チェックポイント、スレッドプール、および消費の同時実行性を設定できます。
注意事項
-
増分ログの保持期間は、テーブルのストリームログの有効期限と同じで、最大 7 日間です。 BaseAndStream トンネルの場合、この期間内に全量データ消費が完了しないと、増分データ消費が開始されるときに
OTSTunnelExpiredエラーが返されます。 このトンネルは、増分データの消費を継続できなくなります。 -
増分データの消費が遅れて保持期間を過ぎた場合、トンネルは利用可能な最新のデータから再開される場合があります。その結果、一部のデータが消費されない可能性があります。
-
期限切れのトンネルは無効になる場合があります。トンネルが無効な状態が 30 日以上続いた場合、そのトンネルは削除され、復元できなくなります。
前提条件
Tablestore SDK for Java をインストールして、TunnelClient を初期化します。
機能の説明
TunnelWorker は、トンネル ID によってトンネルに接続し、ハートビートを使用して現在のクライアントに割り当てられたチャネルを取得し、データを継続的にプルして、レコードの各バッチを IChannelProcessor に渡します。複数の TunnelWorker インスタンスが同じトンネルを消費する場合、サーバーはクライアント間でチャネルを分散します。
トンネルからデータを消費するには、次の手順を実行します:
-
IChannelProcessorを実装します。processメソッドを使用して各バッチを処理し、shutdownメソッドを使用してコールバックで使用されるリソースを解放します。 -
TunnelWorkerConfigオブジェクトを作成して、コールバックと消費動作を設定します。 -
トンネル ID、
TunnelClient、およびTunnelWorkerConfigを使用してTunnelWorkerを作成します。 -
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 |
初期化された |
|
workerConfig (必須) |
TunnelWorkerConfig |
コールバックと消費の動作に関する設定。 |
消費設定
workerConfig は TunnelWorkerConfig 型で、以下のパラメーターが含まれます。
|
名前 |
タイプ |
説明 |
|
channelProcessor (必須) |
IChannelProcessor |
データ処理のコールバック。このパラメーターは、3 つのパラメーターを持つ |
|
heartbeatTimeoutInSec (オプション) |
long |
ハートビートのタイムアウト時間 (秒) です。デフォルト値は |
|
heartbeatIntervalInSec (オプション) |
long |
ハートビートの間隔 (秒) です。デフォルト値は |
|
checkpointIntervalInMillis (オプション) |
long |
サーバー上で消費チェックポイントが記録される間隔 (ミリ秒単位) です。デフォルト値は |
|
clientTag (オプション) |
文字列 |
クライアント ID を生成し、 |
|
readRecordsExecutor (オプション) |
ThreadPoolExecutor |
データをプルするスレッドプールです。デフォルトのプールには、 |
|
processRecordsExecutor (オプション) |
ThreadPoolExecutor |
データを処理するスレッドプール。そのデフォルト設定は |
|
maxChannelParallel (オプション) |
int |
同時にデータをプルおよび処理するチャネルの最大数です。このパラメーターを使用して、メモリ使用量を制限します。デフォルト値は |
|
channelHelperExecutor (オプション) |
ThreadPoolExecutor |
チャネルの初期化、パイプラインのスケジュール、ランタイムエラーの処理を行うヘルパースレッドプール。このパラメーターが設定されていない場合、キャッシュされたスレッドプールが使用されます。 |
|
maxRetryIntervalInMillis (オプション) |
int |
増分データプル中の指数バックオフの最大ベース間隔をミリ秒単位で指定します。デフォルト値は |
|
readMaxTimesPerRound (オプション) |
int |
パイプラインの1ラウンドにおける |
|
readMaxBytesPerRound (オプション) |
int |
1 回のパイプラインラウンドで取得されるデータの最大量 (バイト単位) で、デフォルト値は |
|
enableClosingChannelDetect (オプション) |
ブール値 |
|
同一マシン上で複数の TunnelWorker インスタンスを起動する場合、1 つの TunnelWorkerConfig オブジェクトを再利用して、読み取りおよび処理スレッドプールを共有できます。すべてのワーカーが停止した後、config.shutdown() を 1 回だけ呼び出します。
コールバックデータ
process メソッドは、以下のフィールドを含む ProcessRecordsInput オブジェクトを受け取ります。
|
フィールド |
タイプ |
説明 |
|
records |
|
現在のバッチで取得されたレコード。取得するには |
|
nextToken |
文字列 |
次のバッチ用のトークンです。 |
|
traceId |
文字列 |
現在のプルリクエストのトレース ID。 |
|
channelId |
文字列 |
現在のバッチが属するチャネルの ID です。 |
シナリオ例
消費パラメーターのチューニング
消費のスループットやメモリ使用量が要件を満たさない場合は、ハートビートとチェックポイントの間隔、チャネルの同時実行数、ラウンドごとのプル回数とサイズ、および増分プルのバックオフ間隔を調整します。
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);