サブスクリプション機能
DataHub トピックからデータを消費する場合、アプリケーション障害後に処理を再開できるよう、消費オフセットをお客様側で管理する必要があります。これには進捗の保存と、オフセットを格納するサービスの高可用性の確保が必要であり、アプリケーションの複雑さが増します。これを簡素化するために、DataHub はサーバー側に消費オフセットを格納するサブスクリプションサービスを提供します。いくつかの簡単な設定と最小限のコードで、アプリケーションにとって透過的に動作する高可用なオフセット管理サービスを利用できます。サブスクリプションサービスは柔軟なオフセットのリセット機能も提供し、at-least-once 消費セマンティクスをサポートします。たとえば、特定期間のデータに影響する処理エラーを見つけ、そのデータを再消費したい場合は、オフセットを該当時刻にリセットできます。アプリケーションはこの変更を自動的に検知し、再起動なしでデータの再処理を行います。
サブスクリプションの作成
指定したプロジェクト内のトピックに対してサブスクリプションを作成する権限が、アカウントに付与されていることを確認してください。詳細については、「権限コントロール」をご参照ください。手順は次のとおりです。
トピックページを開き、右上の [+ Subscription] をクリックします。サブスクリプションの詳細を入力し、[Create] をクリックします。
Subscription Application:このサブスクリプションを使用するアプリケーションの名前。
Description:サブスクリプションの詳細な説明。
消費チェックポイントの下にある検索ボタンをクリックして、すべてのシャードの消費状況を表示します。
使用例
サブスクリプション機能はオフセットを格納します。サブスクリプション機能は DataHub の読み取りおよび書き込み機能とは独立していますが (「Java SDK」をご参照ください)、データ読み取り後に消費オフセットを格納する必要がある場合、これらの機能と併用されることが一般的です。
// データを消費し、処理中にオフセットをコミットする例。
public void offset_consumption(int maxRetry) {
String endpoint = "<YourEndPoint>";
String accessId = "<YourAccessId>";
String accessKey = "<YourAccessKey>";
String projectName = "<YourProjectName>";
String topicName = "<YourTopicName>";
String subId = "<YourSubId>";
String shardId = "0";
List<String> shardIds = Arrays.asList(shardId);
// DatahubClient インスタンスを作成します。
DatahubClient datahubClient = DatahubClientBuilder.newBuilder()
.setDatahubConfig(
new DatahubConfig(endpoint,
// バイナリ転送を有効にするか。この機能はサーバーバージョン 2.12 以降でサポートされています。
new AliyunAccount(accessId, accessKey), true))
.build();
RecordSchema schema = datahubClient.getTopic(projectName, topicName).getRecordSchema();
OpenSubscriptionSessionResult openSubscriptionSessionResult = datahubClient.openSubscriptionSession(projectName, topicName, subId, shardIds);
SubscriptionOffset subscriptionOffset = openSubscriptionSessionResult.getOffsets().get(shardId);
// 1. 現在のオフセットに対応するカーソルを取得します。現在のオフセットが期限切れ、または未消費の場合は、ライフサイクル内の先頭レコードのカーソルを取得します。
String cursor = "";
// シーケンス番号が 0 未満の場合、そのシャードは未消費であることを示します。
if (subscriptionOffset.getSequence() < 0) {
// ライフサイクル内の先頭レコードのカーソルを取得します。
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
} else {
// 次のレコードのカーソルを取得します。
long nextSequence = subscriptionOffset.getSequence() + 1;
try {
// SEQUENCE を使用してカーソルを取得すると SeekOutOfRangeException がスローされることがあります。これは、現在のカーソルのデータが期限切れであることを示します。
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
} catch (SeekOutOfRangeException e) {
// ライフサイクル内の先頭レコードのカーソルを取得します。
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
}
}
// 2. レコードを読み取り、オフセットを保存します。この例ではタプルデータを読み取り、1,000 レコードごとにオフセットをコミットします。
long recordCount = 0L;
// 1 回に 1,000 レコードを読み取ります。
int fetchNum = 1000;
int retryNum = 0;
int commitNum = 1000;
while (retryNum < maxRetry) {
try {
GetRecordsResult getRecordsResult = datahubClient.getRecords(projectName, topicName, shardId, schema, cursor, fetchNum);
if (getRecordsResult.getRecordCount() <= 0) {
// データがありません。スリープして再試行します。
System.out.println("no data, sleep 1 second");
Thread.sleep(1000);
continue;
}
for (RecordEntry recordEntry : getRecordsResult.getRecords()) {
// データを処理します。
TupleRecordData data = (TupleRecordData) recordEntry.getRecordData();
System.out.println("field1:" + data.getField("field1") + "\t"
+ "field2:" + data.getField("field2"));
// データ処理後にオフセットを更新します。
recordCount++;
subscriptionOffset.setSequence(recordEntry.getSequence());
subscriptionOffset.setTimestamp(recordEntry.getSystemTime());
// 1,000 レコードごとにオフセットをコミットします。
if (recordCount % commitNum == 0) {
// オフセットをコミットします。
Map<String, SubscriptionOffset> offsetMap = new HashMap<>();
offsetMap.put(shardId, subscriptionOffset);
datahubClient.commitSubscriptionOffset(projectName, topicName, subId, offsetMap);
System.out.println("commit offset successful");
}
}
cursor = getRecordsResult.getNextCursor();
} catch (SubscriptionOfflineException | SubscriptionSessionInvalidException e) {
// 終了します。SubscriptionOfflineException:サブスクリプションがオフラインです。SubscriptionSessionInvalidException:別のクライアントが同じサブスクリプションを消費しています。
e.printStackTrace();
throw e;
} catch (SubscriptionOffsetResetException e) {
// オフセットがリセットされました。最新バージョンの SubscriptionOffset を取得する必要があります。
SubscriptionOffset offset = datahubClient.getSubscriptionOffset(projectName, topicName, subId, shardIds).getOffsets().get(shardId);
subscriptionOffset.setVersionId(offset.getVersionId());
// オフセットのリセット後は、新しいカーソルを取得する必要があります。カーソル取得方法は、オフセットのリセット方法と一致させる必要があります。
// リセット時にシーケンスとタイムスタンプの両方が設定された場合、SEQUENCE または SYSTEM_TIME のいずれかを使用してカーソルを取得できます。
// シーケンスのみが設定された場合は、SEQUENCE を使用する必要があります。
// タイムスタンプのみが設定された場合は、SYSTEM_TIME を使用する必要があります。
// 一般的なルールとして、まず SEQUENCE でカーソルの取得を試み、次に SYSTEM_TIME を試します。両方が失敗した場合は OLDEST を使用します。
cursor = null;
if (cursor == null) {
try {
long nextSequence = offset.getSequence() + 1;
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SEQUENCE, nextSequence).getCursor();
System.out.println("get cursor successful");
} catch (DatahubClientException exception) {
System.out.println("get cursor by SEQUENCE failed, try to get cursor by SYSTEM_TIME");
}
}
if (cursor == null) {
try {
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.SYSTEM_TIME, offset.getTimestamp()).getCursor();
System.out.println("get cursor successful");
} catch (DatahubClientException exception) {
System.out.println("get cursor by SYSTEM_TIME failed, try to get cursor by OLDEST");
}
}
if (cursor == null) {
try {
cursor = datahubClient.getCursor(projectName, topicName, shardId, CursorType.OLDEST).getCursor();
System.out.println("get cursor successful");
} catch (DatahubClientException exception) {
System.out.println("get cursor by OLDEST failed");
System.out.println("get cursor failed!!");
throw e;
}
}
} catch (LimitExceededException e) {
// 上限を超過しました。再試行します。
e.printStackTrace();
retryNum++;
} catch (DatahubClientException e) {
// その他のエラーです。再試行します。
e.printStackTrace();
retryNum++;
} catch (Exception e) {
e.printStackTrace();
System.exit(-1);
}
}
}アプリケーションの初回起動時は、利用可能な最も古いレコードからデータの消費を開始します。アプリケーションの実行中に Web コンソールのサブスクリプションページを更新すると、シャードの消費オフセットが進んでいくことを確認できます。
コンシューマーの実行中に、Web コンソールでチェックポイントのリセット機能を使用してオフセットを手動で変更した場合、アプリケーションは変更を自動的に検知し、新しいオフセットから消費を再開します。これを行うには、クライアントで
SubscriptionOffsetResetExceptionをキャッチし、getSubscriptionOffsetメソッドを呼び出して、最新のSubscriptionOffsetオブジェクトをサーバーから取得します。複数のコンシューマースレッドまたはプロセスで、サブスクリプションの同一シャードを同時に消費しないでください。異なるコンシューマーによってオフセットが上書きされ、格納されたオフセットが未定義の状態になります。この場合、サーバーは
SubscriptionSessionInvalidExceptionをスローします。この例外をキャッチしてアプリケーションを終了し、コンシューマーの重複がないか設計を確認してください。