O Alibaba Cloud Logstash inclui o plug-in logstash-input-jdbc por padrão. Use esse plug-in com uma configuração de pipeline para extrair continuamente dados completos ou incrementais do ApsaraDB RDS for MySQL e indexá-los no Alibaba Cloud Elasticsearch, sem necessidade de instalação adicional. Essa abordagem é adequada para sincronizar todos os dados com latência de alguns segundos ou para consultar dados específicos em um determinado momento.
Como funciona
O plug-in consulta o ApsaraDB RDS for MySQL em um cronograma configurável por meio de uma instrução SQL. A cada execução, ele registra o valor mais alto de uma coluna de rastreamento (como update_time) e usa esse valor como limite inferior para a próxima consulta. Esse mecanismo de polling captura linhas recém-inseridas e atualizadas, mas não detecta linhas excluídas. As exclusões devem ser tratadas separadamente no Elasticsearch.
Limitações
Antes de começar, observe as seguintes restrições:
Mesmo fuso horário obrigatório. A instância RDS, o cluster Elasticsearch e o cluster Logstash devem estar no mesmo fuso horário. Incompatibilidades causam deslocamento de fuso horário nos campos de timestamp após a sincronização.
IDs de documento devem corresponder à chave primária do MySQL. Defina
document_idcomo%{id}na configuração de output, ondeidé a coluna de chave primária no MySQL. Esse mapeamento permite que o plug-in sobrescreva o documento correto no Elasticsearch quando uma linha for atualizada no MySQL. O Elasticsearch trata as sobrescritas da mesma forma que uma operação de atualização padrão: ele exclui o documento antigo e indexa o novo.Exclusões não são propagadas. O plug-in consulta apenas as linhas inseridas ou atualizadas após
sql_last_value, portanto, não detecta operaçõesDELETEno MySQL. Para remover um documento do Elasticsearch, execute o comando de exclusão correspondente diretamente no cluster.Cada linha deve incluir uma coluna de timestamp de inserção/atualização. O plug-in usa essa coluna como referência de rastreamento para determinar quais linhas buscar na próxima execução. Por padrão,
sql_last_valuecomeça em quinta-feira, 1º de janeiro de 1970.Parâmetros de segurança são obrigatórios. Sempre anexe
allowLoadLocalInfile=false&autoDeserialize=falseaojdbc_connection_string. Sem esses parâmetros, a verificação da configuração do pipeline falha.
Pré-requisitos
Antes de iniciar, certifique-se de ter:
Uma instância ApsaraDB RDS for MySQL (este exemplo usa MySQL 5.7)
Um cluster Alibaba Cloud Elasticsearch (este exemplo usa Elasticsearch V7.10)
Um cluster Alibaba Cloud Logstash
Implante todos os três recursos na mesma Virtual Private Cloud (VPC) para manter o tráfego privado e evitar configurações de NAT. Caso sua instância RDS use um endpoint público, configure uma entrada de Source Network Address Translation (SNAT) para o cluster Logstash, ative o endpoint público na instância RDS e adicione os endereços IP dos nós do Logstash à lista de permissões do RDS. Para mais detalhes, consulte Configurar um gateway NAT para transmissão de dados pela Internet.
Sincronizar dados do ApsaraDB RDS for MySQL para o Elasticsearch
Etapa 1: Preparar o ambiente
Ative o recurso Auto Indexing no cluster Elasticsearch para que o Logstash possa criar índices automaticamente. Para mais detalhes, consulte Acessar e configurar um cluster Elasticsearch.
Faça upload de um driver JDBC compatível com sua versão do MySQL para o cluster Logstash. Este exemplo usa
mysql-connector-java-5.1.48.jar. Para mais detalhes, consulte Configurar bibliotecas de terceiros.-
Crie uma tabela de teste e insira dados de amostra na instância RDS:
CREATE table food ( id int PRIMARY key AUTO_INCREMENT, name VARCHAR (32), insert_time DATETIME, update_time DATETIME ); INSERT INTO food values(null,'Chocolates',now(),now()); INSERT INTO food values(null,'Yogurt',now(),now()); INSERT INTO food values(null,'Ham sausage',now(),now()); Adicione os endereços IP dos nós do cluster Logstash à lista de permissões do RDS. Encontre os endereços IP na página Basic Information do cluster Logstash.
Etapa 2: Configurar um pipeline do Logstash
Acesse a página Logstash Clusters no console Alibaba Cloud Elasticsearch.
Na barra de navegação superior, selecione a região onde o cluster está localizado. Na página Logstash Clusters, clique em Pipelines e, em seguida, clique em Create Pipeline.
No painel de navegação à esquerda, clique em Pipelines e, em seguida, clique em Create Pipeline.
-
Insira um Pipeline ID e cole a seguinte configuração no campo Config Settings. Substitua os valores de espaço reservado pelos seus próprios dados.
Para obter o ID do cluster Logstash, consulte Visualizar as informações básicas de um cluster . Para a lista completa de parâmetros de input suportados, consulte Plug-in de input Jdbc do Logstash . Para parâmetros de output, consulte Arquivos de configuração do Logstash .
input { jdbc { # JDBC driver class for MySQL jdbc_driver_class => "com.mysql.jdbc.Driver" # Path to the driver JAR uploaded in Step 1.2. # Format: /ssd/1/share/<Logstash cluster ID>/logstash/current/config/custom/<driver filename> jdbc_driver_library => "/ssd/1/share/<Logstash cluster ID>/logstash/current/config/custom/mysql-connector-java-5.1.48.jar" # RDS connection string. Use the internal endpoint when instances share a VPC. # allowLoadLocalInfile=false and autoDeserialize=false are required — omitting them causes a check failure. jdbc_connection_string => "jdbc:mysql://rm-bp1xxxxx.mysql.rds.aliyuncs.com:3306/<Database name>?useUnicode=true&characterEncoding=utf-8&useSSL=false&allowLoadLocalInfile=false&autoDeserialize=false" # Credentials for the RDS database user jdbc_user => "xxxxx" jdbc_password => "xxxx" # Enable paging and set page size to control memory usage on large tables (default: false) jdbc_paging_enabled => "true" jdbc_page_size => "50000" # Fetch only rows updated after the last recorded value (incremental sync) statement => "select * from food where update_time >= :sql_last_value" # Run every minute (Rufus cron expression) schedule => "* * * * *" # Persist the last run value so the next run starts where this one ended record_last_run => true # Store the tracking value in this file. Use the /ssd/1/<cluster ID>/logstash/data/ path # to ensure write permissions. Replace <Logstash cluster ID> with your cluster ID. last_run_metadata_path => "/ssd/1/<Logstash cluster ID>/logstash/data/last_run_metadata_update_time.txt" # Set clean_run to true to reset and re-sync all data from the beginning clean_run => false # Track by timestamp column to capture inserts and updates. # The values in this column must be sorted in ascending order. tracking_column_type => "timestamp" use_column_value => true tracking_column => "update_time" } } filter { } output { elasticsearch { hosts => "http://es-cn-0h****dd0hcbnl.elasticsearch.aliyuncs.com:9200" index => "rds_es_dxhtest_datetime" user => "elastic" password => "xxxxxxx" # Map the Elasticsearch document ID to the MySQL primary key document_id => "%{id}" } } -
Clique em Next para configurar os parâmetros do pipeline.
AvisoSalvar e implantar o pipeline aciona uma reinicialização do cluster Logstash. Confirme se a reinicialização não afetará suas cargas de trabalho antes de prosseguir.
Parâmetro
Descrição
Padrão
Pipeline Workers
Quantidade de threads que executam plug-ins de filtro e output em paralelo. Aumente este valor quando os recursos de CPU estiverem subutilizados ou houver acúmulo de eventos.
Número de vCPUs
Pipeline Batch Size
Máximo de eventos que um único worker coleta das entradas antes de executar filtros e outputs. Valores maiores aumentam o throughput, mas exigem mais memória heap da JVM.
125
Pipeline Batch Delay
Tempo de espera do worker por eventos adicionais antes de iniciar um pequeno lote, em milissegundos.
50
Queue Type
Modelo de fila interna para buffer de eventos. MEMORY: fila em memória. PERSISTED: fila baseada em disco com ACK para maior durabilidade.
MEMORY
Queue Max Bytes
Tamanho máximo da fila no disco. Deve ser menor que a capacidade de disco disponível.
1024 MB
Queue Checkpoint Writes
Máximo de eventos gravados antes que um checkpoint seja forçado (apenas filas persistentes). Defina como 0 para não haver limite.
1024

Clique em Save and Deploy para reiniciar o cluster Logstash e aplicar a configuração imediatamente. Alternativamente, clique em Save para armazenar a configuração e acionar uma alteração no cluster, mas as configurações ainda não entrarão em vigor. Para implantar posteriormente, acesse a página Pipelines, localize o pipeline e clique em Deploy Now na coluna Actions.
Etapa 3: Verificar a sincronização
Faça login no console Kibana do cluster Elasticsearch. Para mais detalhes, consulte Fazer login no console Kibana.
No canto superior esquerdo, clique em
e escolha .-
Na aba Console, execute o seguinte comando para confirmar que as três linhas foram sincronizadas:
GET rds_es_dxhtest_datetime/_count { "query": {"match_all": {}} }Resposta esperada:
{ "count" : 3, "_shards" : { "total" : 1, "successful" : 1, "skipped" : 0, "failed" : 0 } } -
Atualize e insira linhas no MySQL para testar a sincronização incremental:
UPDATE food SET name='Chocolates',update_time=now() where id = 1; INSERT INTO food values(null,'Egg',now(),now()); -
Após a próxima execução agendada (até um minuto), consulte o Elasticsearch para confirmar as alterações:
-
Pesquise a linha atualizada:
GET rds_es_dxhtest_datetime/_search { "query": { "match": { "name": "Chocolates" } } }
-
Consulte todos os documentos:
GET rds_es_dxhtest_datetime/_search { "query": { "match_all": {} } }
-
Perguntas frequentes
Meu pipeline está travado no estado de inicialização, os dados estão inconsistentes após a sincronização ou a conexão com o banco de dados falhou. O que devo fazer?
Verifique primeiro os logs do cluster Logstash. Acesse o console Logstash e use o recurso Query logs para visualizar os detalhes do erro. Para mais informações, consulte Consultar logs.
Se uma atualização de cluster estiver em andamento quando você aplicar uma correção, pause a atualização primeiro. Consulte Visualizar o progresso de uma tarefa do cluster . Após a correção, o sistema reinicia o cluster e retoma a atualização automaticamente.
A tabela a seguir lista causas comuns e soluções:
|
Causa |
Solução |
|
Endereços IP dos nós do Logstash ausentes na lista de permissões do RDS |
Adicione os endereços IP dos nós do Logstash à lista de permissões do RDS. Para obter os endereços IP, consulte Visualizar as informações básicas de um cluster. Para configuração da lista de permissões, consulte Usar um cliente de banco de dados ou a CLI para conectar-se a uma instância ApsaraDB RDS for MySQL. |
|
Sincronização a partir de MySQL autogerenciado no ECS: IPs privados e portas dos nós ausentes no grupo de segurança do ECS |
Adicione os endereços IP privados e as portas internas dos nós do Logstash ao grupo de segurança do ECS. Consulte Adicionar uma regra de grupo de segurança. |
|
Cluster Elasticsearch em VPC diferente da do cluster Logstash |
Adquira um cluster Elasticsearch na mesma VPC ou configure um gateway NAT para acesso via Internet. Consulte Criar um cluster Alibaba Cloud Elasticsearch e Configurar um gateway NAT para transmissão de dados pela Internet. |
|
Endpoint RDS incorreto ou porta não padrão |
Obtenha o endpoint e a porta corretos no console RDS. Consulte Visualizar e gerenciar endpoints e portas de instâncias. Use o endpoint interno. Para acesso via endpoint público, configure primeiro um gateway NAT. |
|
Auto Indexing desativado no cluster Elasticsearch |
Ative o Auto Indexing. Consulte Configurar o arquivo YML. |
|
Carga do cluster muito alta |
Atualize a configuração do cluster. Consulte Atualizar a configuração de um cluster. Para diagnosticar problemas de carga, ative o X-Pack Monitoring no cluster Logstash. Consulte Ativar o recurso X-Pack Monitoring. |
|
Driver JDBC não enviado |
Faça upload do arquivo do driver. Consulte Configurar bibliotecas de terceiros. |
|
|
Instale o plug-in ou remova o parâmetro |
Para mais orientações sobre solução de problemas, consulte Perguntas frequentes sobre transferência de dados usando Logstash.
Como sincronizo dados de várias tabelas MySQL para índices separados no Elasticsearch?
Defina múltiplos blocos jdbc na seção input, atribua um valor type a cada um e use condições if[type] na seção output para rotear os dados de cada tabela para um índice diferente:
input {
jdbc {
jdbc_driver_class => "com.mysql.jdbc.Driver"
jdbc_driver_library => "/ssd/1/share/<Logstash cluster ID>/logstash/current/config/custom/mysql-connector-java-5.1.48.jar"
jdbc_connection_string => "jdbc:mysql://rm-bp1xxxxx.mysql.rds.aliyuncs.com:3306/<Database name>?useUnicode=true&characterEncoding=utf-8&useSSL=false&allowLoadLocalInfile=false&autoDeserialize=false"
jdbc_user => "xxxxx"
jdbc_password => "xxxx"
jdbc_paging_enabled => "true"
jdbc_page_size => "50000"
statement => "select * from tableA where update_time >= :sql_last_value"
schedule => "* * * * *"
record_last_run => true
last_run_metadata_path => "/ssd/1/<Logstash cluster ID>/logstash/data/last_run_metadata_update_time.txt"
clean_run => false
tracking_column_type => "timestamp"
use_column_value => true
tracking_column => "update_time"
type => "A"
}
jdbc {
jdbc_driver_class => "com.mysql.jdbc.Driver"
jdbc_driver_library => "/ssd/1/share/<Logstash cluster ID>/logstash/current/config/custom/mysql-connector-java-5.1.48.jar"
jdbc_connection_string => "jdbc:mysql://rm-bp1xxxxx.mysql.rds.aliyuncs.com:3306/<Database name>?useUnicode=true&characterEncoding=utf-8&useSSL=false&allowLoadLocalInfile=false&autoDeserialize=false"
jdbc_user => "xxxxx"
jdbc_password => "xxxx"
jdbc_paging_enabled => "true"
jdbc_page_size => "50000"
statement => "select * from tableB where update_time >= :sql_last_value"
schedule => "* * * * *"
record_last_run => true
last_run_metadata_path => "/ssd/1/<Logstash cluster ID>/logstash/data/last_run_metadata_update_time.txt"
clean_run => false
tracking_column_type => "timestamp"
use_column_value => true
tracking_column => "update_time"
type => "B"
}
}
output {
if[type] == "A" {
elasticsearch {
hosts => "http://es-cn-0h****dd0hcbnl.elasticsearch.aliyuncs.com:9200"
index => "rds_es_dxhtest_datetime_A"
user => "elastic"
password => "xxxxxxx"
document_id => "%{id}"
}
}
if[type] == "B" {
elasticsearch {
hosts => "http://es-cn-0h****dd0hcbnl.elasticsearch.aliyuncs.com:9200"
index => "rds_es_dxhtest_datetime_B"
user => "elastic"
password => "xxxxxxx"
document_id => "%{id}"
}
}
}