Todos os produtos
Search
Central de documentação

Realtime Compute for Apache Flink:Ingestão de logs em tempo real

Última atualização: Jun 27, 2026

Este tópico explica como criar rapidamente um job de sincronização de dados do Kafka para o Hologres, permitindo ingerir dados de log em tempo real por meio do console do Realtime Compute for Apache Flink.

Pré-requisitos

  • Verifique se seu usuário RAM ou função RAM possui as permissões necessárias para acessar o console do Realtime Compute for Apache Flink. Para mais informações, consulte Permissões.

  • Crie um workspace do Flink. Para mais informações, consulte Criar um workspace.

  • Armazenamento upstream e downstream

    Nota

    Suas instâncias do ApsaraMQ for Kafka e do Hologres devem estar na mesma região e VPC do seu workspace do Realtime Compute for Apache Flink. Caso contrário, será necessário estabelecer conectividade de rede. Para mais informações, consulte Acessar serviços entre VPCs ou Acessar a internet.

Etapa 1: Configurar uma lista de permissões de IP

Para permitir que o Flink acesse suas instâncias do Kafka e do Hologres, adicione o bloco CIDR do seu workspace do Flink às listas de permissões de IP do Kafka e do Hologres.

  1. Obtenha o bloco CIDR da VPC do seu workspace do Flink.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, escolha More > Workspace Details.

    3. Na caixa de diálogo Workspace Details, localize o CIDR block do vSwitch.

      A caixa de diálogo exibe informações básicas do workspace e uma lista de vSwitches. Na coluna CIDR block, visualize o bloco CIDR de cada zona de disponibilidade. Anote esse bloco CIDR para a próxima etapa.

  2. Adicione o bloco CIDR do workspace do Flink à lista de permissões de IP da sua instância do Kafka.

    É necessário configurar a lista de permissões de IP para um endpoint de VPC. Para obter as etapas detalhadas, consulte Configurar uma lista de permissões de IP. Na caixa de diálogo de edição da lista de permissões, clique em Add Whitelist IP para adicionar o bloco CIDR.

  3. Adicione o bloco CIDR do workspace do Flink à lista de permissões de IP da sua instância do Hologres.

    Faça login na sua instância do Hologres e configure sua lista de permissões de IP. Para obter as etapas detalhadas, consulte Lista de permissões de IP. Na página de configuração da lista de permissões no Security Center do HoloWeb, insira o bloco CIDR no campo IP Address da caixa de diálogo Edit IP Whitelist e clique em OK.

Etapa 2: Preparar dados de teste do Kafka

Utilize o conector Faker do Realtime Compute for Apache Flink para gerar e gravar dados no ApsaraMQ for Kafka. Siga estas etapas no console de desenvolvimento do Realtime Compute for Apache Flink.

  1. No console do ApsaraMQ for Kafka, crie um tópico chamado users.

    Para mais informações, consulte Criar um tópico.

  2. Crie um job para gravar dados no ApsaraMQ for Kafka.

    1. Faça login no console de desenvolvimento do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, clique em Console.

    3. No painel de navegação à esquerda, escolha Development > ETL.

    4. Clique no ícone image e, em seguida, clique em New Stream Draft. Insira um File Name e selecione uma Engine Version.

      O Realtime Compute for Apache Flink também oferece diversos modelos de código e de sincronização de dados. Cada um inclui casos de uso específicos, exemplos de código e instruções. Clique em um modelo para aprender rapidamente os recursos e a sintaxe do Realtime Compute for Apache Flink e implementar sua lógica de negócios. Para mais informações, consulte Modelos de código e Modelos de sincronização de dados.

      Parâmetro

      Descrição

      Exemplo

      File Name

      O nome do job.

      Nota

      O nome do job deve ser exclusivo no workspace atual.

      flink-test

      Engine Version

      A versão do mecanismo Flink para o job atual.

      Selecione uma versão rotulada como Recommended ou Stable para maior confiabilidade e desempenho. Para mais informações sobre versões do mecanismo, consulte Notas de versão e Versões do mecanismo.

      vvr-8.0.8-flink-1.17

    5. Clique em Create.

    6. Escreva as instruções SQL para o job.

      Copie o código a seguir para o editor e modifique os parâmetros conforme seu ambiente.

      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;

      A tabela a seguir descreve os parâmetros que precisam ser modificados.

      Parâmetro

      Valor de exemplo

      Descrição

      properties.bootstrap.servers

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

      Os endpoints dos brokers do ApsaraMQ for Kafka.

      Uma lista separada por vírgulas de endpoints de broker no formato 'host:porta'. Você pode encontrar o endpoint Domain Name da VPC na seção Endpoint Information da página Instance Details.

      topic

      users

      O nome do tópico do ApsaraMQ for Kafka.

  3. Inicie o job.

    1. Na página Development > ETL, clique em Deploy.

    2. Na caixa de diálogo Deploy draft, clique em Confirm.

    3. Configure os recursos para o job. Para mais informações, consulte Configurar recursos para um job.

    4. Na página O&M > Deployments, localize o deployment desejado e clique em Start na coluna Actions. Para mais informações sobre a configuração de inicialização, consulte Iniciar um deployment.

    5. Na página Deployments, monitore as informações de execução e o status do deployment.

      Como a fonte Faker gera um fluxo delimitado, o status do deployment muda para FINISHED cerca de um minuto após a inicialização. Quando o deployment termina, os dados já foram gravados no tópico 'users'. O código a seguir mostra uma amostra dos dados formatados em JSON.

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

Etapa 3: Criar e iniciar um job de sincronização de dados

Flink CDC

  1. Faça login no console de desenvolvimento do Realtime Compute for Apache Flink para criar um job de sincronização de dados.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, clique em Console.

    3. No painel de navegação à esquerda, escolha Development > ETL.

    4. Clique no ícone image e, em seguida, clique em New ETL Draft. Insira um name e selecione uma engine version.

      Parâmetro

      Descrição

      Exemplo

      Name

      O nome do job.

      Nota

      O nome do job deve ser exclusivo no projeto.

      flink-test

      Engine Version

      A versão do mecanismo Flink para o job.

      Selecione uma versão rotulada como Recommended ou Stable. Essas versões oferecem maior confiabilidade e desempenho. Para mais informações sobre versões do mecanismo, consulte Notas de versão e Versões do mecanismo.

      vvr-8.0.8-flink-1.17

    5. Clique em Create.

  2. Escreva um job Flink CDC. Copie o código a seguir para o editor e atualize os parâmetros conforme seu ambiente.

    O job a seguir sincroniza dados de tabela formatados em JSON do tópico users no Kafka para a tabela users no esquema test_schema do banco de dados flink_test_db no Hologres.

    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

    A tabela a seguir descreve os parâmetros a serem modificados.

    Parâmetro

    Exemplo

    Descrição

    properties.bootstrap.servers

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

    Os endereços dos brokers do Kafka.

    O formato é uma lista separada por vírgulas de entradas host:porta. Você pode obter o Domain Name Endpoint para o tipo de rede VPC na seção Network Information da página de instance details.

    topic

    users

    O nome do tópico do Kafka.

    endpoint

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

    O endpoint da instância do Hologres.

    O formato é <ip>:<porta>. Você pode obter o endpoint da VPC na seção Network Information da página de instance details no console do Hologres.

    username

    O nome de usuário e a senha do banco de dados do Hologres. Insira o AccessKey ID e o AccessKey Secret da sua conta Alibaba Cloud.

    Importante

    Para evitar o vazamento do seu par de AccessKey, utilize o gerenciamento de variáveis para especificar seu AccessKey ID e AccessKey Secret. Para mais informações, consulte Gerenciar variáveis.

    password

    dbname

    flink_test_db

    O nome do banco de dados do Hologres.

    source-table

    users

    A tabela de origem. Por padrão, o nome do tópico é utilizado.

    sink-table

    test_schema.users

    A tabela de destino. Especifique a tabela no formato schema.table_name.

  3. Clique em Save.

  4. Na página Development > ETL, clique em Deploy.

  5. Na página O&M > Deployments, localize o deployment desejado e clique em Start na coluna Actions. Para mais informações sobre as configurações de inicialização do job, consulte Iniciar um job.

    Após o início do job, você pode visualizar suas informações de execução e status na página Deployments. Esta página exibe uma lista de deployments com métricas como Status, health score, CPU e memory, além de fornecer ações como Start e Stop.

SQL

  1. Faça login no console de desenvolvimento do Realtime Compute for Apache Flink para criar um job de sincronização de dados.

    1. Faça login no console do Realtime Compute for Apache Flink.

    2. Na coluna Actions do workspace desejado, clique em Console.

    3. No painel de navegação à esquerda, escolha Development > ETL e, em seguida, clique em New.

    4. Clique no ícone image e, em seguida, clique em New Stream Draft. Insira um name e selecione uma engine version.

      Parâmetro

      Descrição

      Exemplo

      Name

      O nome do job.

      Nota

      O nome do job deve ser exclusivo no projeto.

      flink-test

      Engine Version

      A versão do mecanismo Flink para o job.

      Selecione uma versão rotulada como Recommended ou Stable. Essas versões oferecem maior confiabilidade e desempenho. Para mais informações sobre versões do mecanismo, consulte Notas de versão e Versões do mecanismo.

      vvr-8.0.8-flink-1.17

    5. Clique em Create.

  2. Escreva o job SQL. Copie o código a seguir para o editor SQL e atualize os parâmetros conforme seu ambiente.

    Use uma instrução INSERT INTO para sincronizar dados do tópico users no Kafka para a tabela users no banco de dados flink_test_db do Hologres.

    O Hologres fornece otimizações especiais para tipos de dados JSON e JSONB. Utilize uma instrução INSERT INTO para gravar dados JSON aninhados no Hologres.

    Este método exige que você crie primeiro a tabela users no Hologres e, em seguida, execute as seguintes instruções SQL para gravar dados na tabela.

    CREATE TEMPORARY TABLE kafka_users (
      `id` INT NOT NULL,
      `address` STRING, -- The data in this column is nested 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', -- Automatically expand nested columns.
      '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;

    A tabela a seguir descreve os parâmetros a serem modificados.

    Parâmetro

    Exemplo

    Descrição

    properties.bootstrap.servers

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

    Os endereços dos brokers do Kafka.

    O formato é uma lista separada por vírgulas de entradas host:porta. Você pode obter o Domain Name Endpoint para o tipo de rede VPC na seção Network Information da página de instance details.

    topic

    users

    O nome do tópico do Kafka.

    endpoint

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

    O endpoint da instância do Hologres.

    O formato é <ip>:<porta>. Você pode obter o endpoint da VPC na seção Network Information da página de instance details no console do Hologres.

    username

    O nome de usuário e a senha do banco de dados do Hologres. Insira o AccessKey ID e o AccessKey Secret da sua conta Alibaba Cloud.

    Importante

    Para evitar o vazamento do seu par de AccessKey, utilize o gerenciamento de variáveis para especificar seu AccessKey ID e AccessKey Secret. Para mais informações, consulte Gerenciar variáveis.

    password

    dbname

    flink_test_db

    O nome do banco de dados do Hologres.

    tablename

    users

    O nome da tabela do Hologres.

    Nota
    • Se você usar uma instrução INSERT INTO para sincronizar dados, deverá criar a tabela users e suas colunas no banco de dados de destino antecipadamente.

    • Se o esquema não for public, você deve especificar o parâmetro tablename no formato schema.table_name.

  3. Clique em Save.

  4. Na página Development > ETL, clique em Deploy.

  5. Na página O&M > Deployments, localize o deployment desejado e clique em Start na coluna Actions. Para mais informações sobre as configurações de inicialização do job, consulte Iniciar um job.

    Após o início do job, você pode visualizar suas informações de execução e status na página Deployments. Esta página exibe uma lista de deployments com métricas como Status, health score, CPU e memory, além de fornecer ações como Start e Stop.

Etapa 4: Visualizar o resultado da sincronização completa

  1. Faça login no Console de Gerenciamento do Hologres.

  2. Na página Instances, clique no nome da instância desejada.

  3. No canto superior direito da página, clique em Connect to Instance.

  4. Na aba Metadata Management, visualize o esquema da tabela e os dados da tabela users sincronizada no banco de dados flink_test_db.

    Na árvore de navegação à esquerda, expanda o nome da sua instância > flink_test_db > test_schema > Tables. A tabela users sincronizada aparecerá.

    O esquema e os dados da tabela sincronizada são os seguintes.

    • Esquema da tabela

      Clique duas vezes no nome da tabela users para visualizar o esquema da tabela.

      O esquema da tabela users contém os seguintes campos: id (BIGINT, chave primária), first_name (TEXT), last_name (TEXT), address.country (TEXT), address.state (TEXT) e address.city (TEXT).

      Nota

      Durante a sincronização completa de dados, recomendamos definir a partição e o offset de metadados do Kafka como chave primária da tabela do Hologres. Isso evita duplicação de dados caso o job sofra failover e retransmita dados.

    • Dados da tabela

      No canto superior direito da página da tabela users, clique em Query table. Insira o seguinte comando e clique em Run.

      SELECT * FROM test_schema.users;

      O comando retorna o seguinte resultado.

      A consulta retorna várias linhas. Isso confirma que os registros foram sincronizados com sucesso para a tabela users. Cada linha retornada contém dados completos para as colunas id, first_name, last_name, address.country, address.state e address.city.

Etapa 5: Observar a sincronização automática de esquema

  1. No console do ApsaraMQ for Kafka, envie manualmente uma mensagem que contenha uma nova coluna.

    1. Faça login no console do ApsaraMQ for Kafka.

    2. Na página Instances, clique no nome da instância desejada.

    3. Na página Topics, clique no nome do tópico desejado (users).

    4. Clique em Send Message.

    5. Configure a mensagem.

      Na caixa de diálogo Start to Send and Consume Message, configure os parâmetros da seguinte forma.

      Parâmetro

      Exemplo

      Method of Sending

      Selecione Console.

      Message Key

      Insira flinktest.

      Message Content

      Copie e cole o seguinte conteúdo JSON no campo Message Content.

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

      Neste exemplo, house-points é uma nova coluna aninhada.

      Send to Specified Partition

      Selecione Yes.

      Partition ID

      Insira 0.

    6. Clique em OK.

  2. No console do Hologres, visualize as alterações de esquema e dados na tabela users.

    1. Faça login no console do Hologres.

    2. Na página Instances, clique no nome da instância desejada.

    3. No canto superior direito da página, clique em Connect to Instance.

    4. Na aba Metadata Management, clique duas vezes no nome da tabela users.

    5. Clique em Query table, insira a seguinte instrução e, em seguida, clique em Running.

      SELECT * FROM test_schema.users;
    6. Visualize os resultados da consulta.

      A consulta retorna o seguinte resultado.

      O resultado mostra que o registro com ID 100001 foi gravado com sucesso no Hologres, e duas novas colunas, house-points.house e house-points.points, foram adicionadas à tabela do Hologres.

      Nota

      A mensagem enviada ao ApsaraMQ for Kafka contém apenas uma coluna aninhada: house-points. No entanto, como json.infer-schema.flatten-nested-columns.enable está especificado na cláusula WITH, o Realtime Compute for Apache Flink achata automaticamente essa coluna, utilizando os caminhos de acesso aos campos aninhados como nomes das novas colunas.

Referências