Todos os produtos
Search
Central de documentação

:Usar o Kafka para gravar dados no mecanismo de streaming do Lindorm

Última atualização: Jun 28, 2026

O mecanismo de streaming do Lindorm é totalmente compatível com Flink SQL. Use o Flink SQL para criar uma tarefa de computação em tempo real no mecanismo de streaming do Lindorm e processar com eficiência os dados brutos armazenados em um tópico do Apache Kafka. Este tópico descreve como usar o Flink SQL para enviar uma tarefa de computação que importa dados de um tópico do Apache Kafka para uma tabela ampla do Lindorm.

Pré-requisitos

Observações

Se a sua aplicação estiver implantada em uma instância ECS e você precisar acessar a instância do Lindorm por meio de uma VPC, garanta que ambas atendam às seguintes condições de conectividade de rede:

  • Estão na mesma região. Recomendamos que estejam também na mesma zona para reduzir a latência de rede.

  • A instância ECS e a instância do Lindorm pertencem à mesma VPC.

Procedimento

Etapa 1: Preparar os dados

  1. Use a API do Kafka para gravar no tópico os dados que deseja processar. Escolha um dos métodos abaixo para gravar os dados:

    Neste tópico, usamos a ferramenta de script Kafka open source como exemplo para gravar os dados.

    # Create a topic.
    ./kafka-topics.sh --bootstrap-server <Lindorm Stream Kafka Endpoint> --topic log_topic --create
    
    # Write data to the topic.
    ./kafka-console-producer.sh --bootstrap-server <Lindorm Stream Kafka Endpoint> --topic log_topic
    {"loglevel": "INFO", "thread":"thread-1", "class": "com.alibaba.stream.test", "detail":"thread-1 info detail", "timestamp": "1675840911549"}
    {"loglevel": "ERROR", "thread":"thread-2", "class": "com.alibaba.stream.test", "detail":"thread-2 error detail", "timestamp": "1675840911549"}
    {"loglevel": "WARN", "thread":"thread-3", "class": "com.alibaba.stream.test", "detail":"thread-3 warn detail", "timestamp": "1675840911549"}
    {"loglevel": "ERROR", "thread":"thread-4", "class": "com.alibaba.stream.test", "detail":"thread-4 error detail", "timestamp": "1675840911549"}

    Para mais informações sobre como visualizar o Lindorm Stream Kafka endpoint, consulte Visualizar endpoints.

  2. Crie uma tabela de resultados no LindormTable para armazenar a saída do processamento.

    1. Conecte-se ao LindormTable usando o Lindorm-cli. Para mais informações, consulte Usar o Lindorm-cli para conectar-se e utilizar o LindormTable.

    2. Crie uma tabela de resultados chamada log.

      CREATE TABLE IF NOT EXISTS log (
        loglevel VARCHAR,
        thread VARCHAR,
        class VARCHAR,
        detail VARCHAR,
        timestamp BIGINT,
      primary key (loglevel, thread) );

Etapa 2: Instale o cliente do mecanismo de streaming do Lindorm

  1. Execute o comando abaixo na instância ECS para baixar o pacote do cliente do mecanismo de streaming do Lindorm:

    wget https://hbaseuepublic.oss-cn-beijing.aliyuncs.com/lindorm-sqlline-2.0.2.tar.gz
  2. Descompacte o pacote executando o seguinte comando:

    tar zxvf lindorm-sqlline-2.0.2.tar.gz
  3. Acesse o caminho lindorm-sqlline-2.0.2/bin e execute o comando a seguir para se conectar ao mecanismo de streaming do Lindorm:

    ./lindorm-sqlline -url <Lindorm Stream SQL Endpoint>

    Para mais detalhes sobre como visualizar o Lindorm Stream SQL endpoint, consulte Visualizar endpoints.

Etapa 3: Envie uma tarefa de computação no mecanismo de streaming do Lindorm

No exemplo desta etapa, realize as seguintes operações:

  1. Crie um job Flink chamado log_to_lindorm e duas tabelas: originalData e lindorm_log_table. A tabela originalData atua como tabela de origem associada ao tópico do Kafka, enquanto a lindorm_log_table funciona como tabela de destino (sink) para armazenar os logs resultantes.

  2. Defina um job de streaming para filtrar os logs com loglevel igual a ERROR e gravá-los na tabela de resultados.

Código de exemplo:

CREATE FJOB log_to_lindorm(
    --Create the Kafka source table.
    CREATE TABLE originalData(
        `loglevel` VARCHAR,
        `thread` VARCHAR,
        `class` VARCHAR,
        `detail` VARCHAR,
        `timestamp` BIGINT
    )WITH(
        'connector'='kafka',
        'topic'='log_topic',
        'scan.startup.mode'='earliest-offset',
        'properties.bootstrap.servers'='Lindorm Stream Kafka Endpoint',
        'format'='json'
    );
    -- Create the Lindorm wide table.
    CREATE TABLE lindorm_log_table(
        `loglevel` VARCHAR,
        `thread` VARCHAR,
        `class` VARCHAR,
        `detail` VARCHAR,
        `timestamp` BIGINT,
        PRIMARY KEY (`loglevel`, `thread`) NOT ENFORCED
    )WITH(
        'connector'='lindorm',
        'seedServer'='LindormTable Endpoint for HBase APIs',
        'userName'='root',
        'password'='test',
        'tableName'='log',
        'namespace'='default'
    );
    --Filter out the ERROR logs from the data in the Kafka topic and write the logs to the result wide table.
    INSERT INTO lindorm_log_table SELECT * FROM originalData WHERE loglevel = 'ERROR';
);
Nota

Etapa 4: Visualize o resultado do processamento

Escolha um dos métodos abaixo para consultar o resultado do processamento:

  • Conecte-se ao LindormTable via Lindorm-cli e execute o comando seguinte para obter o resultado:

    SELECT * FROM log LIMIT 5;

    O retorno será semelhante ao seguinte:

    +----------+----------+-------------------------+-----------------------+---------------+
    | loglevel |  thread  |          class          |        detail         |   timestamp   |
    +----------+----------+-------------------------+-----------------------+---------------+
    | ERROR    | thread-2 | com.alibaba.stream.test | thread-2 error detail | 1675840911549 |
    | ERROR    | thread-4 | com.alibaba.stream.test | thread-4 error detail | 1675840911549 |
    +----------+----------+-------------------------+-----------------------+---------------+                                     
  • Use o sistema de gerenciamento de cluster do LindormTable para consultar os resultados. Para mais detalhes, consulte Consulta de dados.