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
O mecanismo de streaming está ativado na sua instância do Lindorm. Para mais informações, consulte Ativar o mecanismo de streaming.
Adicione o endereço IP do seu cliente à lista de permissões da instância do Lindorm. Para mais informações, consulte Configurar listas de permissões.
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
-
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.
-
Crie uma tabela de resultados no LindormTable para armazenar a saída do processamento.
Conecte-se ao LindormTable usando o Lindorm-cli. Para mais informações, consulte Usar o Lindorm-cli para conectar-se e utilizar o LindormTable.
-
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
-
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 -
Descompacte o pacote executando o seguinte comando:
tar zxvf lindorm-sqlline-2.0.2.tar.gz -
Acesse o caminho
lindorm-sqlline-2.0.2/bine 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:
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.
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';
);
Para saber como visualizar o endpoint do LindormTable para APIs HBase, consulte Visualizar endpoints.
Para mais informações sobre o conector utilizado na tabela ampla, consulte Configure conectores de tabela ampla para o mecanismo de streaming do Lindorm.
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.