MySQL同步ES的 种方案

发布时间:2026/7/25 17:19:12
MySQL同步ES的 种方案 MySQL同步ES的3种方案在微服务架构和搜索引擎普及的今天MySQL作为关系型数据库ESElasticsearch作为全文搜索引擎两者配合使用是常见的数据存储与检索方案。但MySQL的数据如何实时或准实时同步到ES呢本文从实战角度出发介绍3种主流方案并附上可运行代码示例。## 方案一基于Logstash的准实时同步### 原理Logstash是ELK栈的核心组件通过JDBC输入插件定期轮询MySQL将数据增量或全量同步到ES。这种方式配置简单适合对实时性要求不高的场景如分钟级同步。### 实战代码以下是一个Logstash配置示例实现MySQL表orders到ES索引orders_index的增量同步。ruby# logstash.confinput { jdbc { # MySQL连接配置 jdbc_connection_string jdbc:mysql://localhost:3306/ecommerce?useSSLfalseserverTimezoneUTC jdbc_user root jdbc_password password jdbc_driver_library /usr/share/logstash/mysql-connector-java-8.0.26.jar jdbc_driver_class com.mysql.cj.jdbc.Driver # 全量同步SQL首次运行 statement SELECT * FROM orders WHERE updated_at :sql_last_value # 增量同步字段记录上次同步时间 use_column_value true tracking_column updated_at tracking_column_type timestamp # 轮询间隔秒 schedule */30 * * * * * # 避免重复导入 record_last_run true last_run_metadata_path /usr/share/logstash/.logstash_jdbc_last_run }}filter { # 数据清洗将MySQL的decimal转为float mutate { convert { amount float } } # 添加ES需要的字段 mutate { add_field { [metadata][index_id] %{id} } }}output { elasticsearch { hosts [http://localhost:9200] index orders_index document_id %{id} # 使用data_stream模式支持时序数据 data_stream true } # 调试输出 stdout { codec rubydebug }}### 运行方式bash# 启动Logstashlogstash -f /path/to/logstash.conf优点零编码、配置简单、支持复杂SQL。缺点轮询有延迟、不适合高频数据变更。## 方案二基于Canal的实时同步### 原理阿里巴巴开源的Canal监听MySQL binlog解析为JSON格式推送到Kafka/RocketMQ再由消费者写入ES。实现秒级甚至毫秒级同步。### 实战代码以下是一个Python消费者示例从Kafka读取Canal消息并写入ES。python# canal_consumer.pyimport jsonfrom kafka import KafkaConsumerfrom elasticsearch import Elasticsearchfrom elasticsearch.helpers import bulk# 配置ES连接es Elasticsearch([http://localhost:9200])# Kafka消费者配置consumer KafkaConsumer( canal-topic, # Canal投递的topic bootstrap_servers[localhost:9092], auto_offset_resetlatest, value_deserializerlambda m: json.loads(m.decode(utf-8)))def sync_to_es(message): 将Canal消息同步到ES 消息格式{ type: INSERT/UPDATE/DELETE, data: [{id: 1, name: 张三, age: 25}], database: ecommerce, table: users } actions [] for msg in consumer: canal_data msg.value action_type canal_data[type] database canal_data[database] table canal_data[table] index_name f{database}_{table} # 索引名ecommerce_users for row in canal_data[data]: doc_id row[id] if action_type DELETE: # 删除操作 actions.append({ _op_type: delete, _index: index_name, _id: doc_id }) elif action_type in (INSERT, UPDATE): # 插入或更新操作 actions.append({ _op_type: index, _index: index_name, _id: doc_id, _source: row }) # 批量写入ES if actions: success, _ bulk(es, actions) print(f同步完成: {success}条记录操作类型: {action_type}) actions.clear()if __name__ __main__: print(开始监听Canal消息...) sync_to_es(consumer)### 部署流程1. 启动Canal配置MySQL binlog监听2. 启动Kafka3. 运行Python消费者脚本优点真正实时、支持DDL变更、性能高。缺点架构复杂、依赖组件多、需维护binlog。## 方案三基于应用层双写的强一致性方案### 原理在业务代码中对MySQL执行写操作后同步或异步调用ES API。使用分布式事务如TCC或本地消息表保证最终一致性。### 实战代码以下是一个Spring Boot服务示例使用本地消息表实现可靠同步。java// OrderService.javaServicepublic class OrderService { Autowired private OrderMapper orderMapper; Autowired private EventMapper eventMapper; // 本地消息表 Autowired private RestTemplate restTemplate; Transactional public void createOrder(Order order) { // 1. 写入MySQL orderMapper.insert(order); // 2. 写入本地消息表状态待处理 Event event new Event(); event.setEventType(ORDER_CREATE); event.setPayload(JSON.toJSONString(order)); event.setStatus(0); // 0:待处理, 1:已处理 eventMapper.insert(event); // 3. 同步调用ES失败则通过定时任务重试 try { syncToES(order); event.setStatus(1); eventMapper.updateById(event); } catch (Exception e) { log.error(ES同步失败消息已持久化等待重试, e); // 事务提交后由定时任务处理失败消息 } } private void syncToES(Order order) { String esUrl http://localhost:9200/orders/_doc/ order.getId(); HttpEntityOrder request new HttpEntity(order); restTemplate.exchange(esUrl, HttpMethod.PUT, request, Void.class); } // 定时任务重试失败的消息 Scheduled(fixedDelay 5000) public void retryFailedEvents() { ListEvent failedEvents eventMapper.selectByStatus(0); for (Event event : failedEvents) { Order order JSON.parseObject(event.getPayload(), Order.class); try { syncToES(order); event.setStatus(1); eventMapper.updateById(event); } catch (Exception e) { log.warn(重试失败消息ID: {}, event.getId(), e); } } }}优点强一致性、可控性高、无额外中间件。缺点代码侵入性强、性能有损耗、需处理幂等。## 方案对比总结| 方案 | 实时性 | 复杂度 | 一致性 | 适用场景 ||------|--------|--------|--------|----------|| Logstash | 分钟级 | 低 | 最终一致 | 日志、报表等非关键数据 || Canal | 秒级 | 高 | 最终一致 | 电商、搜索等实时查询 || 应用双写 | 毫秒级 | 中 | 强一致 | 金融、订单等核心业务 |## 总结选择MySQL同步ES的方案需根据业务场景权衡- 如果数据变更不频繁、对实时性要求不高Logstash是最简单的入门选择。- 如果需要秒级同步且能接受架构复杂度Canal是行业标准方案。- 对于核心业务数据如支付、库存应用层双写本地消息表能提供最强一致性保障。实际生产中经常混合使用核心业务用双写消息表非核心用Canal历史数据用Logstash全量导入。无论哪种方案都要注意幂等性设计和失败重试机制避免数据不一致。希望本文的代码示例能为你提供参考帮你在项目中少走弯路。

相关新闻

最新新闻

日新闻

周新闻

月新闻