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

Realtime Compute for Apache Flink:リアルタイムログの取り込み

最終更新日:Sep 19, 2026

このトピックでは、Realtime Compute for Apache Flink コンソールを使用して、Kafka から Hologres へリアルタイムのログデータをインジェストするデータ同期ジョブを手早く作成する方法を説明します。

前提条件

  • RAM ユーザーまたは RAM ロールに、Realtime Compute for Apache Flink コンソールにアクセスするために必要な権限が付与されていること。詳細については、「権限」をご参照ください。

  • Realtime Compute for Apache Flink ワークスペースが作成されていること。詳細については、「ワークスペースの作成」をご参照ください。

  • 上流および下流ストレージ

    説明

    ApsaraMQ for Kafka インスタンスと Hologres インスタンスは、Realtime Compute for Apache Flink ワークスペースと同じリージョンおよび VPC にある必要があります。そうでない場合は、ネットワーク接続を確立する必要があります。詳細については、「VPC をまたいだサービスへのアクセス」または「インターネットへのアクセス」をご参照ください。

ステップ 1:IP ホワイトリストの設定

Flink から Kafka と Hologres のインスタンスにアクセスできるように、Flink ワークスペースの CIDR ブロックを Kafka と Hologres の IP ホワイトリストに追加します。

  1. Flink ワークスペースの VPC CIDR ブロックを取得します。

    1. Realtime Compute for Apache Flink コンソールにログインします。

    2. 対象の[ワークスペース]の [Actions] 列で、[More] > [Workspace Details] を選択します。

    3. [Workspace Details] ダイアログボックスで、VSwitch の [CIDR block] を確認します。

      ダイアログボックスには、ワークスペースの基本情報と VSwitch のリストが表示されます。[CIDR block] 列で、各アベイラビリティーゾーンの CIDR ブロックを確認します。次のステップのために、この CIDR ブロックを控えてください。

  2. Flink ワークスペースの CIDR ブロックを Kafka インスタンスの IP ホワイトリストに追加します。

    VPC エンドポイントの IP ホワイトリストを設定する必要があります。詳細な手順については、「IP ホワイトリストの設定」をご参照ください。ホワイトリスト編集ダイアログボックスで、[Add Whitelist IP] をクリックして CIDR ブロックを追加します。

  3. Flink ワークスペースの CIDR ブロックを Hologres インスタンスの IP ホワイトリストに追加します。

    Hologres インスタンスにログインし、その IP ホワイトリストを設定します。詳細な手順については、「IP ホワイトリスト」をご参照ください。HoloWeb の [Security Center] のホワイトリスト設定ページで、[Edit IP Whitelist] ダイアログボックスの [IP Address] フィールドに CIDR ブロックを入力し、[OK] をクリックします。

ステップ 2:Kafka テストデータの準備

Realtime Compute for Apache Flink の開発コンソールで、Faker コネクタを使用して ApsaraMQ for Kafka にデータを生成して書き込むには、次の手順に従います。

  1. ApsaraMQ for Kafka コンソールで、「users」というトピックを作成します。

    詳細については、「トピックの作成」をご参照ください。

  2. ApsaraMQ for Kafka にデータを書き込むジョブを作成します。

    1. Realtime Compute for Apache Flink の開発コンソールにログオンします。

    2. 対象のワークスペースの [Actions] 列で、[Console] をクリックします。

    3. 左側メニューで、[Development] > [ETL] を選択します。

    4. image アイコンをクリックし、[New Blank Stream Draft] をクリックします。[Name] を入力し、[Engine Version] を選択します。

      Realtime Compute for Apache Flink には、さまざまなコードおよびデータ同期テンプレートも用意されています。各テンプレートには、特定のユースケース、コードサンプル、手順が含まれています。テンプレートをクリックすると、Realtime Compute for Apache Flink の機能と構文をすばやく習得し、ビジネスロジックを実装できます。詳細については、「コードテンプレート」および「データ同期テンプレート」をご参照ください。

      パラメータ

      説明

      例

      [Name]

      ジョブの名前。

      説明

      ジョブ名は、現在のワークスペース内で一意である必要があります。

      flink-test

      [Engine Version]

      現在のジョブの Flink エンジンバージョン。

      信頼性とパフォーマンスを向上させるには、[Recommended] または [Stable] のラベルが付いたバージョンを選択してください。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。

      vvr-8.0.8-flink-1.17

    5. [Create] をクリックします。

    6. ジョブの SQL ステートメントを記述します。

      次のコードをエディタにコピーし、環境に合わせてパラメータを変更してください。

      CREATE TEMPORARY TABLE source (
        id INT,
        first_name STRING,
        last_name STRING,
        `address` ROW<`country` STRING, `state` STRING, `city` STRING>,
        event_time TIMESTAMP
      ) WITH (
        'connector' = 'faker',
        'number-of-rows' = '100',
        'rows-per-second' = '10',
        'fields.id.expression' = '#{number.numberBetween ''0'',''1000''}',
        'fields.first_name.expression' = '#{name.firstName}',
        'fields.last_name.expression' = '#{name.lastName}',
        'fields.address.country.expression' = '#{address.country}',
        'fields.address.state.expression' = '#{address.state}',
        'fields.address.city.expression' = '#{address.city}',
        'fields.event_time.expression' = '#{date.past ''15'',''SECONDS''}'
      );
      
      CREATE TEMPORARY TABLE sink (
        id INT,
        first_name STRING,
        last_name STRING,
        `address` ROW<`country` STRING, `state` STRING, `city` STRING>,
        `timestamp` TIMESTAMP METADATA
      ) WITH (
        'connector' = 'kafka',
        'properties.bootstrap.servers' = 'alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092',
        'topic' = 'users',
        'format' = 'json',
        'properties.enable.idempotence'='false'
      );
      
      INSERT INTO sink SELECT id, first_name, last_name, `address` FROM source;

      次の表に、変更する必要があるパラメータを示します。

      パラメータ

      値の例

      説明

      properties.bootstrap.servers

      alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092

      ApsaraMQ for Kafka ブローカーエンドポイント。

      カンマで区切られた、「host:port」形式のブローカーエンドポイントのリスト。VPC の [Domain Name] エンドポイントは、[Instance Details] ページの [Endpoint Information] セクションで確認できます。

      topic

      users

      ApsaraMQ for Kafka トピック名。

  3. ジョブを起動します。

    1. [Development] > [ETL] ページで、[Deploy] をクリックします。

    2. [Deploy draft] ダイアログボックスで、[Confirm] をクリックします。

    3. ジョブのリソースを設定します。詳細については、「ジョブのリソースの設定」をご参照ください。

    4. [O&M] > [Deployments] ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。起動設定の詳細については、「デプロイメントの起動」をご参照ください。

    5. [Deployments] ページで、デプロイメントのランタイム情報とステータスを監視できます。

      Faker ソースは有界ストリームを生成するため、デプロイメントのステータスは起動後約 1 分で FINISHED に変わります。デプロイメントが完了すると、データが「users」トピックに書き込まれています。以下は、JSON 形式のデータサンプルの例です。

      {
        "id": 765,
        "first_name": "Barry",
        "last_name": "Pollich",
        "address": {
          "country": "United Arab Emirates",
          "state": "Nevada",
          "city": "Powlowskifurt"
        }
      }

ステップ 3:データ同期ジョブの作成と開始

Flink CDC

  1. Realtime Compute for Apache Flink 開発コンソールにログインして、データ同期ジョブを作成します。

    1. Realtime Compute for Apache Flink コンソール にログインします。

    2. 対象のワークスペースの [Actions] 列で、[Console] をクリックします。

    3. 左側のナビゲーションペインで、[Development] > [データインジェスト] を選択します。

    4. image アイコンをクリックし、[新規下書き] をクリックします。[name] を入力し、[engine version] を選択します。

      パラメータ

      説明

      例

      [名前]

      ジョブの名前。

      説明

      ジョブ名はプロジェクト内で一意である必要があります。

      flink-test

      [エンジンバージョン]

      ジョブの Flink エンジンバージョン。

      [Recommended] または [Stable] のラベルが付いたバージョンを選択します。これらのバージョンは、より高い信頼性とパフォーマンスを提供します。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。

      vvr-8.0.8-flink-1.17

    5. [Create] をクリックします。

  2. Flink CDC ジョブを記述します。次のコードをエディターにコピーし、ご使用の環境に合わせてパラメータを更新します。

    次のジョブは、Kafka の users トピックから JSON 形式のテーブルデータを、Hologres の flink_test_db データベースの test_schema スキーマにある users テーブルに同期します。

    source:
      type: kafka
      name: Kafka Source
      properties.bootstrap.servers: alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092
      topic: users
      scan.startup.mode: earliest-offset
      value.format: json
      json.infer-schema.flatten-nested-columns.enable: true
    
    sink:
      type: hologres
      name: Hologres Sink
      endpoint: hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80
      dbname: flink_test_db
      username: ******
      password: **
      sink.type-normalize-strategy: ONLY_BIGINT_OR_TEXT
    
    transform:
      - source-table: \.*.\.*
        projection: \*
        primary-keys: id
        
    route:
      - source-table: users
        sink-table: test_schema.users

    次の表に、変更する必要があるパラメータを示します。

    パラメータ

    例

    説明

    properties.bootstrap.servers

    alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092

    Kafka ブローカーのアドレス。

    形式は、host:port 形式のエントリをカンマで区切ったリストです。[インスタンス詳細] ページの [Network Information] セクションから、VPC ネットワークタイプの [Domain Name Endpoint] を取得できます。

    topic

    users

    Kafka トピックの名前。

    endpoint

    hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80

    Hologres インスタンスのエンドポイント。

    形式は <ip>:<port> です。Hologres コンソールの [インスタンス詳細] ページの [Network Information] セクションから VPC エンドポイントを取得できます。

    username

    **

    Hologres データベースのユーザー名とパスワード。Alibaba Cloud アカウントの AccessKey ID と AccessKey Secret を入力します。

    重要

    AccessKey ペアの漏洩を防ぐために、変数管理を使用して AccessKey ID と AccessKey Secret を指定することを推奨します。詳細については、「変数の管理」をご参照ください。

    password

    **

    dbname

    flink_test_db

    Hologres データベースの名前。

    source-table

    users

    ソーステーブル。デフォルトでは、トピック名が使用されます。

    sink-table

    test_schema.users

    宛先テーブル。テーブルを schema.table_name 形式で指定します。

  3. [Save] をクリックします。

  4. [Development] > [データインジェスト] ページで、[Deploy] をクリックします。

  5. [運用保守] > [Deployments] ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。ジョブの起動設定の詳細については、「ジョブの開始」をご参照ください。

    ジョブが開始されると、[Deployments] ページでそのランタイム情報とステータスを表示できます。このページには、[Status]、[ヘルススコア]、[CPU]、[memory] などのメトリクスを含むデプロイメントのリストが表示され、[Start] や [Stop] などのアクションを実行できます。

SQL

  1. Realtime Compute for Apache Flink 開発コンソールにログインして、データ同期ジョブを作成します。

    1. Realtime Compute for Apache Flink コンソール にログインします。

    2. 対象のワークスペースの [Actions] 列で、[Console] をクリックします。

    3. 左側のナビゲーションペインで、[開発] > [ETL] を選択し、次に [新規] をクリックします。

    4. image アイコンをクリックし、次に [新規ブランクストリームドラフト] をクリックします。[名前] を入力し、[エンジンバージョン] を選択します。

      パラメータ

      説明

      例

      [名前]

      ジョブの名前。

      説明

      ジョブ名はプロジェクト内で一意である必要があります。

      flink-test

      [エンジンバージョン]

      ジョブの Flink エンジンバージョン。

      [Recommended] または [Stable] のラベルが付いたバージョンを選択します。これらのバージョンは、より高い信頼性とパフォーマンスを提供します。エンジンバージョンの詳細については、「リリースノート」および「エンジンバージョン」をご参照ください。

      vvr-8.0.8-flink-1.17

    5. [Create] をクリックします。

  2. SQL ジョブを記述します。次のコードを SQL エディターにコピーし、ご使用の環境に合わせてパラメータを更新します。

    INSERT INTO ステートメントを使用して、Kafka の users トピックから Hologres の flink_test_db データベースの users テーブルにデータを同期できます。

    Hologres は、JSON および JSONB データ型に対して特別な最適化を提供します。INSERT INTO ステートメントを使用して、ネストされた JSON データを Hologres に書き込むことができます。

    この方法では、まず Hologres に users テーブルを作成し、次に以下の SQL ステートメントを実行してテーブルにデータを書き込む必要があります。

    CREATE TEMPORARY TABLE kafka_users (
      `id` INT NOT NULL,
      `address` STRING, -- この列のデータはネストされた JSON です。
      `offset` BIGINT NOT NULL METADATA,
      `partition` BIGINT NOT NULL METADATA,
      `timestamp` TIMESTAMP METADATA,
      `date` AS CAST(`timestamp` AS DATE),
      `country` AS JSON_VALUE(`address`, '$.country')
    ) WITH (
      'connector' = 'kafka',
      'properties.bootstrap.servers' = 'alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092',
      'topic' = 'users',
      'format' = 'json',
      'json.infer-schema.flatten-nested-columns.enable' = 'true', -- ネストされた列を自動的に展開します。
      'scan.startup.mode' = 'earliest-offset'
    );
    
    CREATE TEMPORARY TABLE holo (
      `id` INT NOT NULL,
      `address` STRING,
      `offset` BIGINT,
      `partition` BIGINT,
      `timestamp` TIMESTAMP,
      `date` DATE,
      `country` STRING
    ) WITH (
      'connector' = 'hologres',
      'endpoint' = 'hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80',
      'username' = '******',
      'password' = '******',
      'dbname' = 'flink_test_db',
      'tablename' = 'users'
    );
    
    INSERT INTO holo
    SELECT * FROM kafka_users;

    次の表に、変更する必要があるパラメータを示します。

    パラメータ

    例

    説明

    properties.bootstrap.servers

    alikafka-host1.aliyuncs.com:9092,alikafka-host2.aliyuncs.com:9092,alikafka-host3.aliyuncs.com:9092

    Kafka ブローカーのアドレス。

    形式は、host:port 形式のエントリをカンマで区切ったリストです。[インスタンス詳細] ページの [Network Information] セクションから、VPC ネットワークタイプの [Domain Name Endpoint] を取得できます。

    topic

    users

    Kafka トピックの名前。

    endpoint

    hgpostcn-****-cn-beijing-vpc.hologres.aliyuncs.com:80

    Hologres インスタンスのエンドポイント。

    形式は <ip>:<port> です。Hologres コンソールの [インスタンス詳細] ページの [Network Information] セクションから VPC エンドポイントを取得できます。

    username

    ******

    Hologres データベースのユーザー名とパスワード。Alibaba Cloud アカウントの AccessKey ID と AccessKey Secret を入力します。

    重要

    AccessKey ペアの漏洩を防ぐために、変数管理を使用して AccessKey ID と AccessKey Secret を指定することを推奨します。詳細については、「変数の管理」をご参照ください。

    password

    ******

    dbname

    flink_test_db

    Hologres データベースの名前。

    tablename

    users

    Hologres テーブルの名前。

    説明
    • INSERT INTO ステートメントを使用してデータを同期する場合、宛先データベースに users テーブルとその列を事前に作成する必要があります。

    • スキーマが public でない場合は、tablename パラメータを schema.table_name 形式で指定する必要があります。

  3. [Save] をクリックします。

  4. [Development] > [ETL] ページで、[Deploy] をクリックします。

  5. [運用保守] > [Deployments] ページで、対象のデプロイメントを見つけ、[Actions] 列の [Start] をクリックします。ジョブの起動設定の詳細については、「ジョブの開始」をご参照ください。

    ジョブが開始されると、[Deployments] ページでそのランタイム情報とステータスを表示できます。このページには、[Status]、[ヘルススコア]、[CPU]、[memory] などのメトリクスを含むデプロイメントのリストが表示され、[Start] や [Stop] などのアクションを実行できます。

ステップ 4: 完全同期結果の表示

  1. Hologres 管理コンソールにログインします。

  2. [Instances] ページで、ターゲットインスタンスの名前をクリックします。

  3. ページの右上隅で、[Connect to Instance] をクリックします。

  4. [Metadata Management] タブで、flink_test_db データベース内の同期された users テーブルのテーブルスキーマとデータを表示します。

    左側のナビゲーションウィンドウで、インスタンス名 > [flink_test_db] > [test_schema] > [Tables] を展開します。同期された [users] テーブルが表示されます。

    同期されたテーブルのスキーマとデータは次のとおりです。

    • テーブルスキーマ

      users テーブル名をダブルクリックして、テーブルスキーマを表示します。

      users テーブルスキーマには、次のフィールドが含まれます: id (BIGINT、プライマリキー)、first_name (TEXT)、last_name (TEXT)、address.country (TEXT)、address.state (TEXT)、address.city (TEXT)。

      説明

      完全データ同期中は、Kafka メタデータのパーティションとオフセットを Hologres テーブルのプライマリキーとして定義することを推奨します。これにより、ジョブがフェイルオーバーしてデータを再送信する際のデータ重複を防ぐことができます。

    • テーブルデータ

      users テーブルページの右上隅で、[Query table] をクリックします。次のコマンドを入力し、[Run] をクリックします。

      SELECT * FROM test_schema.users;

      コマンドは次の結果を返します:

      クエリは複数の行を返します。これにより、レコードが users テーブルに正常に同期されたことが確認できます。返された各行には、id、first_name、last_name、address.country、address.state、address.city 列の完全なデータが含まれています。

ステップ 5:自動スキーマ同期の確認

  1. ApsaraMQ for Kafka コンソールで、新しい列を含むメッセージを手動で送信します。

    1. ApsaraMQ for Kafka コンソールにログインします。

    2. [Instances] ページで、対象インスタンスの名前をクリックします。

    3. [Topics] ページで、対象トピック (users) の名前をクリックします。

    4. [Send Message] をクリックします。

    5. メッセージを設定します。

      [Start to Send and Consume Message] ダイアログボックスで、パラメーターを次のように設定します:

      パラメーター

      例

      [送信方法]

      [Console] を選択します。

      [メッセージキー]

      「flinktest」と入力します。

      [メッセージ内容]

      以下の JSON コンテンツをコピーして、「メッセージ内容」フィールドに貼り付けてください。

      {
        "id": 100001,
        "first_name": "Ichiro",
        "last_name": "Tanaka",
        "address": {
          "country": "Isle of Man",
          "state": "Montana",
          "city": "East Coleburgh"
        },
        "house-points": {
          "house": "Pukwudgie",
          "points": 76
        }
      }
      説明

      この例では、 house-points は新しいネストされた列です。

      [指定パーティションへの送信]

      [Yes] を選択します。

      [パーティション ID]

      「0」と入力します。

    6. [OK] をクリックします。

  2. Hologres コンソールで、users テーブルのスキーマとデータの変更を確認します。

    1. Hologres コンソールにログインします。

    2. [Instances] ページで、対象インスタンスの名前をクリックします。

    3. ページの右上隅にある [インスタンスに接続] をクリックします。

    4. [メタデータ管理] タブで、users テーブル名をダブルクリックします。

    5. [テーブルのクエリ] をクリックし、次の文を入力して [Run] をクリックします。

      SELECT * FROM test_schema.users;
    6. クエリ結果を表示します。

      クエリは次の結果を返します:

      この結果から、ID 100001 のレコードが Hologres に正常に書き込まれ、2 つの新しい列 (house-points.house と house-points.points) が Hologres テーブルに追加されたことがわかります。

      説明

      ApsaraMQ for Kafka に送信されたメッセージには、ネストされた列が 1 つ (house-points) しか含まれていません。しかし、WITH 句で json.infer-schema.flatten-nested-columns.enable が指定されているため、Realtime Compute for Apache Flink はこの列を自動的にフラット化し、ネストされたフィールドアクセスパスを新しい列名として使用します。

参考資料

  • Message Queue for Apache Kafka をソーステーブルまたは結果テーブルとして使用する方法については、「Message Queue for Apache Kafka」をご参照ください。

  • ノードの並列度とリソースを調整してジョブのパフォーマンスを向上させるには、「ジョブデプロイメントの設定」をご参照ください。