DataHub コネクタを使用すると、Alibaba Cloud DataHub から Flink ジョブにストリーミングデータを読み込み、処理結果を DataHub トピックに書き戻すことができます。Flink SQL と DataStream API の両方をサポートしています。
機能
| 項目 | 説明 |
|---|---|
| サポートタイプ | ソースおよび sink |
| 実行モード | ストリーミングおよびバッチ |
| データフォーマット | N/A |
| メトリクス | N/A |
| API タイプ | DataStream および SQL |
| Sink でのデータ更新/削除のサポート | サポートされていません。Sink は挿入のみの行をターゲットトピックに書き込みます。 |
前提条件
開始する前に、以下が準備できていることを確認してください:
-
DataHub プロジェクトとトピック。詳細については、「DataHub の概要」をご参照ください。
-
DataHub サブスクリプション (ソースに必須)。詳細については、「サブスクリプションの作成」をご参照ください。
-
Alibaba Cloud の AccessKey ID と AccessKey Secret。詳細については、「コンソール操作」をご参照ください。
制限事項
-
バッチジョブで DataHub ソーステーブルを読み取ることは推奨されません。バッチモードでは、DataHub ソーステーブルは終了状態に達することができず、ジョブは完了せずに待機し続けます。
構文
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
コネクタオプション
一般
| オプション | タイプ | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
| connector | String | はい | (なし) | コネクタタイプ。これを datahub に設定します。 |
| endPoint | String | はい | (なし) | DataHub プロジェクトのエンドポイント。値はリージョンによって異なります。詳細については、「エンドポイント」をご参照ください。 |
| project | String | はい | (なし) | DataHub プロジェクト名。 |
| topic | String | はい | (なし) | DataHub トピック名。BLOB トピック (型なし、非構造化データ) の場合、Flink テーブルには VARBINARY 列が 1 つだけ含まれている必要があります。 |
| accessId | String | はい | (なし) | Alibaba Cloud アカウントの AccessKey ID です。 ハードコーディングするのではなく、変数として保存してください。 詳細については、「変数の管理」をご参照ください。 |
| accessKey | String | はい | (なし) | ご利用の Alibaba Cloud アカウントの AccessKey Secret。 |
| retryTimeout | Integer | いいえ | 1800000 | リトライ試行の最大タイムアウト (ミリ秒)。 |
| retryInterval | Integer | いいえ | 1000 | リトライ試行の間隔 (ミリ秒)。 |
| CompressType | String | いいえ | lz4 | 読み取りおよび書き込みの圧縮アルゴリズム。有効な値:lz4、deflate、"" (無効)。VVR 6.0.5 以降が必要です。 |
ソース固有
| オプション | タイプ | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
| subId | String | はい | (なし) | DataHub サブスクリプション ID。 |
| maxFetchSize | Integer | いいえ | 50 | リクエストごとにフェッチされるレコード数。この値を増やすと、読み取りスループットが向上します。 |
| maxBufferSize | Integer | いいえ | 50 | 非同期読み取りからキャッシュされるレコードの最大数。この値を増やすと、読み取りスループットが向上します。 |
| fetchLatestDelay | Integer | いいえ | 500 | データが利用できない場合のスリープ時間 (ミリ秒)。この値を小さくすると、低トラフィックのトピックでの読み取りレイテンシが短縮されます。 |
| lengthCheck | String | いいえ | NONE | 解析されたフィールド数が定義された列数と一致しない行を処理するためのルール。有効な値:NONE、SKIP、EXCEPTION、PAD。詳細については、「フィールド数検証ルール」をご参照ください。 |
| columnErrorDebug | Boolean | いいえ | false | フィールド解析エラーのデバッグロギングを有効にするかどうかを指定します。true に設定すると、解析例外ログが出力されます。 |
| startTime | String | いいえ | (なし) | 消費を開始するタイムスタンプ。フォーマット:yyyy-MM-dd hh:mm:ss。 |
| endTime | String | いいえ | (なし) | 消費を停止するタイムスタンプ。フォーマット:yyyy-MM-dd hh:mm:ss。 |
| startTimeMs | Long | いいえ | -1 | 消費を開始するタイムスタンプ (ミリ秒)。startTime よりも優先されます。詳細については、「消費開始位置」をご参照ください。 |
消費開始位置
startTimeMs オプションは、ソースが読み取りを開始する位置を制御します:
-
-1 (デフォルト):トピックの最新オフセットから開始します。オフセットが存在しない場合は、最も古いオフセットにフォールバックします。
-
特定のタイムスタンプ:指定されたタイムスタンプ以降の最初のレコードから開始します。
デフォルト値の -1 はデータ損失を引き起こす可能性があります。最初のチェックポイントの前にジョブが失敗した場合、トピックの最新オフセットが進んでいる可能性があり、そのウィンドウ中に書き込まれたレコードはスキップされます。開始位置を制御するには、startTimeMs を特定のタイムスタンプに明示的に設定してください。
フィールド数検証ルール
lengthCheck オプションは、行内の解析済みフィールド数が定義された列数と一致しない場合の動作を決定します:
| 値 | 動作 |
|---|---|
NONE (デフォルト) |
解析済みフィールド > 定義済み列の場合:定義された数まで左から右に読み取ります。解析済みフィールド < 定義済み列の場合:行をスキップします。 |
SKIP |
解析済みフィールド数が定義された列数と異なる行をスキップします。 |
EXCEPTION |
解析済みフィールド数が定義された列数と異なる場合に例外をスローします。 |
PAD |
左から右に読み取ります。解析済みフィールド > 定義済み列の場合:定義された数まで左から右に読み取ります。解析済みフィールド < 定義済み列の場合:欠落しているフィールドを null で埋めます。 |
Sink 固有
| オプション | タイプ | 必須 | デフォルト | 説明 |
|---|---|---|---|---|
| batchCount | Integer | いいえ | 500 | 書き込みバッチあたりの最大行数。 |
| batchSize | Integer | いいえ | 512000 | 書き込みバッチの最大サイズ (バイト)。 |
| flushInterval | Integer | いいえ | 5000 | フラッシュ間隔 (ミリ秒)。 |
| hashFields | String | いいえ | null | 行をシャードにルーティングするために使用される列名のカンマ区切りリスト。これらの列で同じ値を持つ行は、同じシャードに書き込まれます。デフォルト (null) ではランダム書き込みが使用されます。例:hashFields=a,b。 |
| timeZone | String | いいえ | (なし) | TIMESTAMP フィールドを変換する際に使用されるタイムゾーン。 |
| schemaVersion | Integer | いいえ | -1 | 登録されたスキーマレジストリ内のスキーマバージョン。 |
書き込みバッチのフラッシュ動作
書き込みバッチは、以下のいずれかの条件が最初に満たされると DataHub にフラッシュされます:
-
バッファリングされた行数が
batchCountに達した。 -
バッファリングされたデータの合計サイズが
batchSizeに達した。 -
最後のフラッシュからの時間が
flushIntervalを超えた。
batchCount、batchSize、または flushInterval を増やすと、レイテンシが高くなる代わりに書き込みスループットが向上します。
データ型マッピング
| Flink 型 | DataHub 型 |
|---|---|
| TINYINT | TINYINT |
| BOOLEAN | BOOLEAN |
| INTEGER | INTEGER |
| BIGINT | BIGINT |
| BIGINT | TIMESTAMP |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DECIMAL | DECIMAL |
| VARCHAR | STRING |
| SMALLINT | SMALLINT |
| VARBINARY | BLOB |
メタデータ
メタデータフィールドは読み取り専用 (R) です。ソーステーブル定義で METADATA VIRTUAL として宣言することで、DataHub に書き戻すことなくクエリに含めることができます。
| キー | データ型 | 説明 | R/W |
|---|---|---|---|
| shard-id | BIGINT METADATA VIRTUAL | レコードのシャード ID。 | R |
| sequence | STRING METADATA VIRTUAL | シャード内でのレコードのシーケンス番号。 | R |
| system-time | TIMESTAMP METADATA VIRTUAL | DataHub がレコードを受信した時刻。 | R |
例
ソース
次の例では、DataHub トピックからデータを読み取り、コンソールに出力します。
CREATE TEMPORARY TABLE datahub_input (
`time` BIGINT,
`sequence` STRING METADATA VIRTUAL,
`shard-id` BIGINT METADATA VIRTUAL,
`system-time` TIMESTAMP METADATA VIRTUAL
) WITH (
'connector' = 'datahub',
'subId' = '<yourSubId>',
'endPoint' = '<yourEndPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}'
);
CREATE TEMPORARY TABLE test_out (
`time` BIGINT,
`sequence` STRING,
`shard-id` BIGINT,
`system-time` TIMESTAMP
) WITH (
'connector' = 'print',
'logger' = 'true'
);
INSERT INTO test_out
SELECT
`time`,
`sequence`,
`shard-id`,
`system-time`
FROM datahub_input;
Sink
次の例では、ある DataHub トピックから読み取り、name フィールドを小文字に変換し、結果を別の DataHub トピックに書き込みます。
CREATE TEMPORARY TABLE datahub_source (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'subId' = '<yourSubId>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'startTime' = '2018-06-01 00:00:00'
);
CREATE TEMPORARY TABLE datahub_sink (
name VARCHAR
) WITH (
'connector' = 'datahub',
'endPoint' = '<endPoint>',
'project' = '<yourProjectName>',
'topic' = '<yourTopicName>',
'accessId' = '${secret_values.ak_id}',
'accessKey' = '${secret_values.ak_secret}',
'batchSize' = '512000',
'batchCount' = '500'
);
INSERT INTO datahub_sink
SELECT
LOWER(name)
FROM datahub_source;
DataStream API
DataHub で DataStream API を使用するには、Realtime Compute for Apache Flink 用の DataStream コネクタを設定します。詳細については、「DataStream コネクタの設定」をご参照ください。
DataHub からの読み取り
VVR は、Flink の SourceFunction インターフェイスを実装する DatahubSourceFunction クラスを提供します。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1);
// DataHub ソースを設定
DatahubSourceFunction datahubSource =
new DatahubSourceFunction(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<yourStartTime>,
<yourEndTime>
);
datahubSource.setRequestTimeout(30 * 1000);
datahubSource.enableExitAfterReadFinished();
env.addSource(datahubSource)
.map((MapFunction<RecordEntry, Tuple2<String, Long>>) this::getStringLongTuple2)
.print();
env.execute();
private Tuple2<String, Long> getStringLongTuple2(RecordEntry recordEntry) {
Tuple2<String, Long> tuple2 = new Tuple2<>();
TupleRecordData recordData = (TupleRecordData) (recordEntry.getRecordData());
tuple2.f0 = (String) recordData.getField(0);
tuple2.f1 = (Long) recordData.getField(1);
return tuple2;
}
DataHub への書き込み
VVR は、DatahubSinkFunction インターフェイスを実装する OutputFormatSinkFunction クラスを提供します。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// DataHub sink を設定
env.generateSequence(0, 100)
.map((MapFunction<Long, RecordEntry>) aLong -> getRecordEntry(aLong, "default:"))
.addSink(
new DatahubSinkFunction<>(
<yourEndPoint>,
<yourProjectName>,
<yourTopicName>,
<yourSubId>,
<yourAccessId>,
<yourAccessKey>,
"public",
<schemaVersion> // スキーマレジストリが有効な場合は、スキーマバージョンを指定する必要があります。それ以外の場合は、これを 0 に設定します。
)
);
env.execute();
private RecordEntry getRecordEntry(Long message, String s) {
RecordSchema recordSchema = new RecordSchema();
recordSchema.addField(new Field("f1", FieldType.STRING));
recordSchema.addField(new Field("f2", FieldType.BIGINT));
recordSchema.addField(new Field("f3", FieldType.DOUBLE));
recordSchema.addField(new Field("f4", FieldType.BOOLEAN));
recordSchema.addField(new Field("f5", FieldType.TIMESTAMP));
recordSchema.addField(new Field("f6", FieldType.DECIMAL));
RecordEntry recordEntry = new RecordEntry();
TupleRecordData recordData = new TupleRecordData(recordSchema);
recordData.setField(0, s + message);
recordData.setField(1, message);
recordEntry.setRecordData(recordData);
return recordEntry;
}
Maven 依存関係
DataHub DataStream コネクタをプロジェクトに追加します。利用可能なすべてのバージョンは、Maven セントラルリポジトリにリストされています。
<dependency>
<groupId>com.alibaba.ververica</groupId>
<artifactId>ververica-connector-datahub</artifactId>
<version>${vvr-version}</version>
</dependency>
次のステップ
-
Realtime Compute for Apache Flink でサポートされているコネクタの完全なリストについては、「サポートされているコネクタ」をご参照ください。
-
代わりに Kafka コネクタを使用して DataHub に接続するには、「Message Queue for Apache Kafka」をご参照ください。