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

Realtime Compute for Apache Flink:DataHub コネクタ

最終更新日:Jul 23, 2026

DataHub コネクタを使用すると、Alibaba Cloud DataHub から Flink ジョブにストリーミングデータを読み込み、処理結果を DataHub トピックに書き戻すことができます。Flink SQL と DataStream API の両方をサポートしています。

説明 DataHub は Kafka プロトコルと互換性があります。Kafka プロトコルを使用して Flink を DataHub に接続するには、Upsert Kafka コネクタではなく、標準の Kafka コネクタを使用します。詳細については、「Kafka との互換性」をご参照ください。

機能

項目 説明
サポートタイプ ソースおよび 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 に書き戻すことなくクエリに含めることができます。

説明 メタデータフィールドは、VVR 3.0.1 以降を使用している場合にのみ利用可能です。
キー データ型 説明 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>

次のステップ