全部產品
Search
文件中心

Realtime Compute for Apache Flink:日誌即時入倉

更新時間:Sep 19, 2026

本文為您介紹如何通過Realtime Compute控制台快速構建從Kafka到Hologres的資料同步作業,實現日誌資料的即時入倉部署。

前提條件

步驟一:配置IP白名單

為了讓Flink能訪問Kafka和Hologres執行個體,您需要將Flink工作空間的網段添加到Kafka和Hologres的白名單中。

  1. 擷取Flink工作空間的VPC網段。

    1. 登入Realtime Compute控制台

    2. 在目標工作空間右側操作列,選擇更多 > 工作空間詳情

    3. 工作空間詳情對話方塊,查看虛擬交換器的網段資訊。

      對話方塊中展示工作空間基本資料及虛擬交換器列表,在列表的網段列可查看各可用性區域對應的VPC網段,將此網段記錄下來用於後續白名單配置。

  2. 在訊息佇列Kafka的IP白名單中,添加Flink工作空間的網段資訊。

    您需要為網路類型為VPC的存取點配置白名單,操作步驟請參見配置白名單。在對應的白名單編輯對話方塊中,單擊添加白名單IP添加網段。

  3. 在Hologres的IP白名單中,添加Flink工作空間的網段資訊。

    登入Hologres執行個體後配置IP白名單,操作步驟請參見IP白名單。在HoloWeb資訊安全中心的白名單配置頁面,在編輯IP白名單對話方塊的IP地址欄填入網段資訊並單擊確認

步驟二:準備Kafka測試資料

使用Realtime ComputeFlink版的類比資料產生Faker作為資料產生器,將資料寫入到Kafka中。請按以下步驟使用Realtime Compute開發控制台將資料寫入至訊息佇列Kafka。

  1. 在Kafka控制台建立一個名稱為users的Topic。

    操作詳情請參見步驟一:建立Topic

  2. 建立將資料寫入到Kafka的作業。

    1. 登入Realtime Compute管理主控台

    2. 單擊目標工作空間操作列下的控制台

    3. 在左側導覽列,單擊資料開發 > ETL

    4. 單擊image後,單擊建立流作業,填寫檔案名稱並選擇引擎版本

      Flink也為您提供了豐富的代碼模板和資料同步,每種代碼模板都為您提供了具體的使用情境、程式碼範例和使用指導。您可以直接單擊對應的模板快速地瞭解Flink產品功能和相關文法,實現您的商務邏輯,詳情請參見代碼模板資料同步模板

      作業參數

      說明

      樣本

      檔案名稱

      作業的名稱。

      說明

      作業名稱在當前專案中必須保持唯一。

      flink-test

      引擎版本

      當前作業使用的Flink引擎版本。

      建議使用帶有推薦穩定標籤的版本,這些版本具有更高的可靠性和效能表現,引擎版本詳情請參見功能發布記錄引擎版本介紹

      vvr-8.0.8-flink-1.17

    5. 單擊建立

    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 * FROM source;

      需要修改的參數配置資訊如下:

      參數

      樣本值

      說明

      properties.bootstrap.servers

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

      Kafka Broker地址。

      格式為host:port,host:port,host:port,以英文逗號(,)分割。您可以在執行個體詳情頁面的存取點資訊地區擷取網路類型為VPC的網域名稱存取點作為該參數的值。

      topic

      users

      Kafka Topic名稱。

  3. 啟動作業。

    1. 資料開發 > ETL頁面,單擊部署

    2. 部署新版本對話方塊中,單擊確定

    3. 配置作業資源,資源設定填寫詳情請參見配置作業資源

    4. 營運中心 > 作業營運頁面,單擊目標作業名稱操作列中的啟動。關於作業啟動的配置說明,請參見作業啟動

    5. 您可以在作業營運頁面觀察作業的運行資訊和狀態。

      由於faker資料來源是一個有限流,因此在作業處於運行狀態後,大約1分鐘左右後,作業就會處於完成狀態。當作業結束運行代表作業已經將相關的資料寫入到Kafka的users中。其中,寫入到訊息佇列Kafka的JSON資料格式大致如下。

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

步驟三:建立並啟動資料同步作業

通過Flink CDC同步

  1. 登入Realtime Compute開發控制台,建立資料同步作業。

    1. 登入Realtime Compute管理主控台

    2. 單擊目標工作空間操作列下的控制台

    3. 在左側導覽列,單擊資料開發 > 資料攝入

    4. 單擊image後,單擊建立資料攝入草稿,填寫檔案名稱並選擇引擎版本

      作業參數

      說明

      樣本

      檔案名稱

      作業的名稱。

      說明

      作業名稱在當前專案中必須保持唯一。

      flink-test

      引擎版本

      當前作業使用的Flink引擎版本。

      建議使用帶有推薦穩定標籤的版本,這些版本具有更高的可靠性和效能表現,引擎版本詳情請參見功能發布記錄引擎版本介紹

      vvr-8.0.8-flink-1.17

    5. 單擊建立

  2. 編寫Flink CDC作業。將以下作業代碼拷貝到作業文本編輯區,然後根據實際配置,修改參數配置資訊。

    假設Kafka的topic 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 Broker地址。

    格式為host:port,host:port,host:port,以英文逗號(,)分割。您可以在執行個體詳情頁面的存取點資訊地區擷取網路類型為VPC的網域名稱存取點作為該參數的值。

    topic

    users

    Kafka Topic名稱。

    endpoint

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

    Hologres端點。

    格式為<ip>:<port>。您可以在Hologres執行個體詳情頁面擷取網路類型為指定VPC的網域名稱資訊作為該參數的值。

    username

    **

    Hologres使用者名稱和密碼,請填寫阿里雲帳號的AccessKey ID和AccessKey Secret。

    重要

    為了避免您的AK資訊泄露,建議您通過密鑰管理的方式填寫AccessKey ID取值,詳情請參見變數管理

    password

    **

    dbname

    flink_test_db

    Hologres資料庫名稱。

    source-table

    users

    定義來自於上遊哪個表,預設是topic名稱。

    sink-table

    test_schema.users

    定義寫入下遊哪張表,使用逗號串連Schema和表名。

  3. 單擊儲存

  4. 資料開發 > 資料攝入頁面,單擊部署

  5. 營運中心 > 作業營運頁面,單擊目標作業名稱操作列中的啟動關於作業啟動的配置說明,請參見作業啟動

    作業啟動後,您可以在作業營運介面觀察作業的運行資訊和狀態。頁面中顯示作業列表,包含狀態健康分CPU記憶體等運行指標,以及啟動停止等操作按鈕。

通過SQL同步

  1. 登入Realtime Compute開發控制台,建立資料同步作業。

    1. 登入Realtime Compute管理主控台

    2. 單擊目標工作空間操作列下的控制台

    3. 在左側導覽列,單擊資料開發 > ETL,單擊建立

    4. 單擊image後,單擊建立流作業,填寫檔案名稱並選擇引擎版本

      作業參數

      說明

      樣本

      檔案名稱

      作業的名稱。

      說明

      作業名稱在當前專案中必須保持唯一。

      flink-test

      引擎版本

      當前作業使用的Flink引擎版本。

      建議使用帶有推薦穩定標籤的版本,這些版本具有更高的可靠性和效能表現,引擎版本詳情請參見功能發布記錄引擎版本介紹

      vvr-8.0.8-flink-1.17

    5. 單擊建立

  2. 編寫SQL作業。將以下作業代碼拷貝到作業文本編輯區,然後根據實際配置,修改參數配置資訊。

    將訊息佇列Kafka中名稱為users的Topic資料同步至Hologres的flink_test_db資料庫的users表中。您可以通過以下INSERT INTO方式完成資料同步。

    考慮到Hologres中對於JSON和JSONB類型的資料會進行特殊的最佳化,您也可以通過INSERT INTO語句將嵌套JSON寫入到Hologres中。

    該方式需要您手動在Hologres中建立users表,然後通過下文的SQL將資料寫入到Hologres的表中。

    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 Broker地址。

    格式為host:port,host:port,host:port,以英文逗號(,)分割。您可以在執行個體詳情頁面的存取點資訊地區擷取網路類型為VPC的網域名稱存取點作為該參數的值。

    topic

    users

    Kafka Topic名稱。

    endpoint

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

    Hologres端點。

    格式為<ip>:<port>。您可以在Hologres執行個體詳情頁面擷取網路類型為指定VPC的網域名稱資訊作為該參數的值。

    username

    ******

    Hologres使用者名稱和密碼,請填寫阿里雲帳號的AccessKey ID和AccessKey Secret。

    重要

    為了避免您的AK資訊泄露,建議您通過密鑰管理的方式填寫AccessKey ID取值,詳情請參見變數管理

    password

    ******

    dbname

    flink_test_db

    Hologres資料庫名稱。

    tablename

    users

    Hologres表名稱。

    說明
    • 如果您通過INSERT INTO方式同步資料,則需要提前在目標執行個體的資料庫中建立users表和欄位。

    • 如果Schema不為Public時,則tablename需要填寫為schema.tableName。

  3. 單擊儲存

  4. 資料開發 > ETL頁面,單擊部署

  5. 營運中心 > 作業營運頁面,單擊目標作業名稱操作列中的啟動。關於作業啟動的配置說明,請參見作業啟動

    作業啟動後,您可以在作業營運介面觀察作業的運行資訊和狀態。頁面中顯示作業列表,包含狀態健康分CPU記憶體等運行指標,以及啟動停止等操作按鈕。

步驟四:觀察全量同步結果

  1. 登入Hologres管理主控台

  2. 執行個體列表頁面,單擊目標執行個體名稱。

  3. 在頁面右上方,單擊登入執行個體

  4. 中繼資料管理頁簽,查看users資料庫中同步的users表結構和資料。

    在左側導航樹中,依次展開目標執行個體名稱 > flink_test_db > test_schema > ,可以看到已同步的users表。

    同步後的表結構和資料如下所示。

    • 表結構

      雙擊users表名稱,查看錶結構。

      users表結構包含以下欄位:id(int8,主鍵)、first_name(text)、last_name(text)、address.country(text)、address.state(text)、address.city(text)。

      說明

      在同步過程中,建議聲明Kafka的Metadata partition和offset作為Hologres表中的主鍵。這樣可以避免由於作業Failover,資料重發導致下遊儲存多份相同資料。

    • 表資料

      在users表資訊頁面右上方,單擊查詢表後,輸入如下命令,單擊運行

      SELECT * FROM test_schema.users;

      表資料結果如下所示。

      查詢結果顯示users表中已成功同步多行資料,包含idfirst_namelast_nameaddress.countryaddress.stateaddress.city列的完整記錄。

步驟五:觀察自動同步表結構變更

  1. 在Kafka控制台手動發送一條包含新增列的訊息。

    1. 登入雲訊息佇列 Kafka 版控制台

    2. 執行個體列表頁面,單擊目標執行個體名稱。

    3. Topic管理頁面,單擊目標Topic名稱users。

    4. 單擊體驗發送訊息

    5. 填寫訊息內容。

      快速體驗訊息收發對話方塊中,參照以下配置填寫各項參數。

      配置項

      樣本

      發送方式

      選中控制台

      訊息Key

      填寫為flinktest。

      訊息內容

      將以下JSON內容複寫粘貼到訊息內容中。

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

      該樣本中house-points是一個新增的嵌套列。

      發送到指定分區

      選中

      分區ID

      填寫為0。

    6. 單擊確定

  2. 在Hologres控制台,查看users表結構和資料的變化。

    1. 登入Hologres管理主控台

    2. 執行個體列表頁面,單擊目標執行個體名稱。

    3. 在頁面右上方,單擊登入執行個體

    4. 中繼資料管理頁簽,雙擊users表名稱。

    5. 單擊查詢表後,輸入如下命令,單擊運行

      SELECT * FROM test_schema.users;
    6. 查看錶資料結果。

      表資料結果如下所示。

      可以觀察到id為100001的資料已經成功地寫入到了Hologres中。同時,Hologres中多了house-points.house和house-points.points 兩列。

      說明

      雖然插入到Kafka中的資料只有一個嵌套列house-points,但是由於在users表的WITH參數內聲明要求json.infer-schema.flatten-nested-columns.enable,那麼Flink 就會自動展平新增的嵌套列,並用訪問該列的路徑作為展開後的列的名字。

相關文檔