Sincronize dados em tempo real do ApsaraMQ for Kafka para o ApsaraDB for ClickHouse usando o mecanismo de tabela Kafka integrado e uma materialized view.
Limitações
A sincronização de dados é compatível apenas com instâncias do ApsaraMQ for Kafka e clusters Kafka autogerenciados implantados em instâncias ECS.
Pré-requisitos
-
ApsaraDB for ClickHouse:
Crie um cluster de destino na mesma região e VPC da instância do ApsaraMQ for Kafka. Para mais informações, consulte Create a cluster.
Crie uma conta de banco de dados com as permissões necessárias para o cluster de destino. Para mais informações, consulte Account Management.
-
ApsaraMQ for Kafka:
Crie um topic has been created.
Crie um consumer group has been created.
Observações de uso
O tópico assinado pela tabela externa Kafka do ApsaraDB for ClickHouse não deve ter outros consumers.
Ao criar a tabela externa Kafka, a materialized view e a tabela local, garanta que os tipos de campo das três tabelas correspondam.
Procedimento
O exemplo a seguir sincroniza dados do ApsaraMQ for Kafka para a tabela distribuída kafka_table_distributed no banco de dados padrão de um cluster Community-compatible Edition do ApsaraDB for ClickHouse.
Etapa 1: Entender o funcionamento da sincronização
O ApsaraDB for ClickHouse usa o mecanismo de tabela Kafka e uma materialized view para consumir e armazenar dados do Kafka em tempo real. O fluxo de dados ocorre da seguinte forma:
Tópico Kafka: contém os dados de origem a sincronizar.
Tabela externa Kafka do ApsaraDB for ClickHouse (usa o mecanismo Kafka): extrai dados de origem de um tópico Kafka específico.
Materialized view: lê dados de origem da tabela externa Kafka e os insere em uma tabela local no ApsaraDB for ClickHouse.
Tabela local: armazena os dados sincronizados.
Etapa 2: Conectar-se ao cluster do ApsaraDB for ClickHouse
Para mais informações, consulte Connect to an ApsaraDB for ClickHouse cluster by using DMS.
Etapa 3: Criar uma tabela externa Kafka
A tabela externa Kafka usa o mecanismo de tabela Kafka para extrair dados de um tópico Kafka especificado. Esta tabela tem as seguintes características:
Por padrão, não é possível consultar diretamente a tabela externa Kafka.
Essa tabela serve exclusivamente para consumir dados do Kafka e não armazena dados. Use uma materialized view para processar e inserir os dados em uma tabela de destino.
A sintaxe para criação da tabela é a seguinte:
Os tipos de campo da tabela externa Kafka devem corresponder aos tipos de dados das mensagens no Kafka.
CREATE TABLE [IF NOT EXISTS] [db.]table_name [ON CLUSTER cluster]
(
name1 [type1] [DEFAULT|MATERIALIZED|ALIAS expr1],
name2 [type2] [DEFAULT|MATERIALIZED|ALIAS expr2],
...
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'host:port1,host:port2,host:port3',
kafka_topic_list = 'topic_name1,topic_name2,...',
kafka_group_name = 'group_name',
kafka_format = 'data_format'[,]
[kafka_row_delimiter = 'delimiter_symbol',]
[kafka_num_consumers = N,]
[kafka_thread_per_consumer = 1,]
[kafka_max_block_size = 0,]
[kafka_skip_broken_messages = N,]
[kafka_commit_every_batch = 0,]
[kafka_auto_offset_reset = N]
A tabela a seguir descreve os parâmetros comuns.
|
Parâmetro |
Obrigatório |
Descrição |
|
kafka_broker_list |
Sim |
Lista de endpoints de broker separados por vírgulas para o cluster Kafka. Para visualizar endpoints, consulte View Endpoints.
|
|
kafka_topic_list |
Sim |
Lista de nomes de tópicos separados por vírgulas. Para visualizar nomes de tópicos, consulte Create a topic. |
|
kafka_group_name |
Sim |
Nome do grupo de consumidores do Kafka. Para mais informações, consulte Create a group. |
|
kafka_format |
Sim |
Formato do corpo da mensagem compatível com o ApsaraDB for ClickHouse. Nota
Para ver os formatos de corpo de mensagem compatíveis com o ApsaraDB for ClickHouse, consulte Formats for Input and Output Data. |
|
kafka_row_delimiter |
Não |
Delimitador usado para separar linhas. O valor padrão é \n. Defina este parâmetro para corresponder ao delimitador real dos seus dados. |
|
kafka_num_consumers |
Não |
Número de consumers por tabela. O valor padrão é 1. Nota
|
|
kafka_thread_per_consumer |
Não |
Define se uma thread dedicado será ativada para cada consumer. O valor padrão é 0. Valores válidos:
Para melhorar a velocidade de consumo, consulte Kafka performance tuning. |
|
kafka_max_block_size |
Não |
Tamanho máximo, em bytes, de um lote de mensagens do Kafka. O valor padrão é 65536. |
|
kafka_skip_broken_messages |
Não |
Número de erros de análise a ignorar. O valor padrão é 0. Se definir |
|
kafka_commit_every_batch |
Não |
Frequência dos commits do Kafka. O valor padrão é 0. Valores válidos:
|
|
kafka_auto_offset_reset |
Não |
Offset inicial para leitura dos dados do Kafka. Valores válidos:
Nota
Clusters do ApsaraDB for ClickHouse com versão de kernel 21.8 não suportam este parâmetro. |
Para mais informações sobre os parâmetros, consulte Kafka.
O código a seguir fornece um exemplo:
CREATE TABLE default.kafka_src_table ON CLUSTER `default`
(
-- Define the fields of the table schema.
id Int32,
name String
) ENGINE = Kafka()
SETTINGS
kafka_broker_list = 'alikafka-post-cn-****-1-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****1-2-vpc.alikafka.aliyuncs.com:9092,alikafka-post-cn-****-3-vpc.alikafka.aliyuncs.com:9092',
kafka_topic_list = 'testforCK',
kafka_group_name = 'GroupForTestCK',
kafka_format = 'CSV';
Etapa 4: Criar uma tabela de destino
Escolha a instrução de criação de tabela correspondente à edição do seu cluster.
Em clusters Enterprise Edition, basta criar uma tabela local. Em clusters Community-compatible Edition, talvez seja necessário criar uma tabela distribuída, dependendo do ambiente e dos requisitos. Os códigos a seguir apresentam exemplos de instruções. Para mais informações sobre a sintaxe de criação de tabelas, consulte CREATE TABLE.
Enterprise edition
CREATE TABLE default.kafka_table_local ON CLUSTER default (
id Int32,
name String
) ENGINE = MergeTree()
ORDER BY (id);
Se ocorrer o erro ON CLUSTER is not allowed for Replicated database ao executar esta instrução, atualize a versão do kernel para resolver o problema. Para atualizar a versão do kernel, consulte Upgrade the minor engine version.
Community-compatible edition
Os mecanismos de tabela diferem entre clusters de réplica única e réplica dupla. Selecione o mecanismo adequado conforme o tipo de réplica do seu cluster.
Ao criar uma tabela em um cluster de réplica dupla, use obrigatoriamente um mecanismo Replicated da família MergeTree. Tabelas com mecanismos não replicados nesse tipo de cluster impedem a replicação de dados entre réplicas, o que pode causar inconsistência de dados.
Single-replica
-
Crie uma tabela local.
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = MergeTree() ORDER BY (id); -
(Opcional) Crie uma tabela distribuída.
Se precisar importar dados apenas para a tabela local, pule esta etapa.
Em clusters com múltiplos nós, recomendamos criar uma tabela distribuída.
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
Double-replica
-
Crie uma tabela local.
CREATE TABLE default.kafka_table_local ON CLUSTER default ( id Int32, name String ) ENGINE = ReplicatedMergeTree() ORDER BY (id); -
(Opcional) Crie uma tabela distribuída.
Se precisar importar dados apenas para a tabela local, pule esta etapa.
Em clusters com múltiplos nós, recomendamos criar uma tabela distribuída.
CREATE TABLE kafka_table_distributed ON CLUSTER default AS default.kafka_table_local ENGINE = Distributed(default, default, kafka_table_local, id);
Etapa 5: Criar uma materialized view
O ApsaraDB for ClickHouse usa uma materialized view para ler dados de origem da tabela externa Kafka e inseri-los em uma tabela local no próprio ApsaraDB for ClickHouse.
A sintaxe para criação de uma materialized view é a seguinte:
Garanta que os campos SELECT correspondam à estrutura da tabela de destino ou use funções de conversão para adequar o formato dos dados à estrutura da tabela de destino.
CREATE MATERIALIZED VIEW <view_name> ON CLUSTER default TO <dest_table> AS SELECT * FROM <src_table>;
A tabela a seguir descreve os parâmetros.
|
Parâmetro |
Obrigatório |
Descrição |
Exemplo |
|
view_name |
Sim |
Nome da view. |
consumer |
|
dest_table |
Sim |
Tabela de destino para armazenar dados do Kafka.
|
|
|
src_table |
Sim |
Tabela externa Kafka. |
kafka_src_table |
O código a seguir fornece exemplos de instruções.
Enterprise edition
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_local AS SELECT * FROM kafka_src_table;
Community-compatible edition
Neste exemplo, os dados de origem são armazenados na tabela distribuída kafka_table_distributed.
CREATE MATERIALIZED VIEW consumer ON CLUSTER default TO kafka_table_distributed AS SELECT * FROM kafka_src_table;
Etapa 6: Verificar a sincronização
-
Envie mensagens para o tópico na instância do ApsaraMQ for Kafka.
Faça login no console do ApsaraMQ for Kafka.
Na página Instance list, clique em nome da instância de destino.
Na página Topics, localize o tópico de destino e selecione na coluna Actions.
-
Na página Send and Consume Message with Quick Experience, insira o Message Content.
Este exemplo envia as mensagens
1,ae2,b. Clique em OK.
-
Conecte-se ao cluster do ApsaraDB for ClickHouse, consulte a tabela distribuída e verifique se os dados foram sincronizados.
Para fazer login em um cluster do ApsaraDB for ClickHouse, consulte Connect to an ApsaraDB for ClickHouse cluster by using DMS.
Use as seguintes instruções para consultar e verificar os dados:
Enterprise edition
SELECT * FROM kafka_table_local;Community-compatible edition
Exemplo de consulta a uma tabela distribuída:
Se a tabela de destino for local, substitua o nome da tabela distribuída na consulta pelo nome da tabela local.
Em um cluster Community-compatible Edition com múltiplos nós, consulte a tabela distribuída. Consultar diretamente uma tabela local retorna dados de apenas um nó, resultando em um conjunto de resultados incompleto.
SELECT * FROM kafka_table_distributed;Se a consulta retornar resultados, a sincronização de dados do Kafka para o ApsaraDB for ClickHouse foi bem-sucedida.
Os resultados da consulta são os seguintes.
┌─id─┬─name─┐ │ 1 │ a │ │ 2 │ b │ └────┴──────┘Se os resultados da consulta não corresponderem ao esperado, prossiga para Step 7 (Optional): Check the consumption status of the Kafka external table para solucionar o problema.
Etapa 7 (Opcional): Verificar o status de consumo do Kafka
Se os dados sincronizados não corresponderem aos dados no Kafka, consulte a tabela do sistema para verificar o status de consumo da tabela externa Kafka e solucionar exceções.
Engine v23.8 or later
Execute a seguinte instrução para consultar a tabela de sistema system.kafka_consumers e visualizar o status de consumo da tabela externa Kafka:
select * from system.kafka_consumers;
A tabela a seguir descreve os campos da tabela system.kafka_consumers.
|
Campo |
Descrição |
|
database |
Banco de dados onde está localizada a tabela externa Kafka. |
|
table |
Nome da tabela externa Kafka. |
|
consumer_id |
ID do consumer do Kafka. Uma tabela pode ter vários consumers. O parâmetro kafka_num_consumers define essa quantidade durante a criação da tabela externa Kafka. |
|
assignments.topic |
Tópico do Kafka. |
|
assignments.partition_id |
ID da partição do Kafka. Cada partição pode ser atribuída a apenas um consumer. |
|
assignments.current_offset |
Offset atual. |
|
exceptions.time |
Timestamps das 10 exceções mais recentes. |
|
exceptions.text |
Texto das 10 exceções mais recentes. |
|
last_poll_time |
Timestamp da última sondagem (polling). |
|
num_messages_read |
Número de mensagens lidas pelo consumer. |
|
last_commit_time |
Timestamp do último commit. |
|
num_commits |
Total de commits executados pelo consumer. |
|
last_rebalance_time |
Timestamp do último rebalanceamento do Kafka. |
|
num_rebalance_revocations |
Número de vezes que partições foram revogadas do consumer. |
|
num_rebalance_assignments |
Número de vezes que partições foram atribuídas ao consumer no cluster Kafka. |
|
is_currently_used |
Indica se o consumer está em uso. |
|
last_used |
Momento em que o consumer foi usado pela última vez, em tempo Unix (microssegundos). |
|
rdkafka_stat |
Estatísticas internas da biblioteca. Para mais informações, consulte librdkafka. O valor padrão é 3000, indicando geração de estatísticas a cada 3 segundos. Nota
Quando |
Engine earlier than v23.8
Execute a seguinte instrução para consultar a tabela de sistema system.kafka e visualizar o status de consumo da tabela externa Kafka:
SELECT * FROM system.kafka;
A tabela a seguir descreve os campos da tabela system.kafka.
|
Campo |
Descrição |
|
database |
Nome do banco de dados onde está localizada a tabela externa Kafka. |
|
table |
Nome da tabela externa Kafka. |
|
topic |
Nome do tópico consumido pela tabela externa Kafka. |
|
consumer_group |
Nome do grupo de consumidores usado pela tabela externa Kafka. |
|
last_read_message_count |
Número de mensagens extraídas da tabela externa Kafka. |
|
status |
Status do consumo de mensagens do Kafka pela tabela externa. Valores válidos:
|
|
exception |
Detalhes da exceção. Nota
Se o valor de status for error, este parâmetro retorna detalhes da exceção. |