Todos os produtos
Search
Central de documentação

Elasticsearch:Usar o Logstash para sincronizar dados do ApsaraDB RDS for MySQL para o Elasticsearch

Última atualização: Jun 27, 2026

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_id como %{id} na configuração de output, onde id é 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ções DELETE no 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_value começa em quinta-feira, 1º de janeiro de 1970.

  • Parâmetros de segurança são obrigatórios. Sempre anexe allowLoadLocalInfile=false&autoDeserialize=false ao jdbc_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

  1. 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.

  2. 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.

  3. 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());
  4. 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

  1. Acesse a página Logstash Clusters no console Alibaba Cloud Elasticsearch.

  2. 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.

  3. No painel de navegação à esquerda, clique em Pipelines e, em seguida, clique em Create Pipeline.

  4. 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}"
      }
    }
  5. Clique em Next para configurar os parâmetros do pipeline.

    Aviso

    Salvar 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

    Pipeline parameter configuration

  6. 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

  1. Faça login no console Kibana do cluster Elasticsearch. Para mais detalhes, consulte Fazer login no console Kibana.

  2. No canto superior esquerdo, clique em 菜单.png e escolha Management > Dev Tools.

  3. 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
      }
    }
  4. 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());
  5. 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.

file_extend especificado na configuração do pipeline, mas logstash-output-file_extend não instalado

Instale o plug-in ou remova o parâmetro file_extend. Consulte Instalar ou remover um plug-in do Logstash.

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}"
        }
    }
}