ARTICLE DETAIL

资讯详情

深耕商务建站与企业官网运营的一线实战洞察。

Canal 与 Elasticsearch 实时同步:增量索引更新、删除处理与全量重建方案

Canal 与 Elasticsearch 实时同步:增量索引更新、删除处理与全量重建方案 1. Canal 与 Elasticsearch 实时同步概述Canal 是阿里巴巴开源的基于数据库增量日志解析的组件支持 MySQL、Oracle 等数据库。它通过解析数据库的 binlog 日志将数据变更事件推送到消息队列或直接处理实现数据的准实时同步。Elasticsearch 是一个基于 Lucene 库的搜索引擎具有强大的全文检索和分析能力。将 Canal 与 Elasticsearch 结合使用可以实现数据库数据到 Elasticsearch 的准实时同步充分发挥 Elasticsearch 在搜索、分析和日志处理方面的优势。这种架构广泛应用于电商搜索、日志分析、监控告警等场景。Canal 与 Elasticsearch 同步的基本架构包括Canal Server: 监听并解析数据库 binlogCanal Client: 接收变更事件并处理数据转换: 将数据库数据转换为 Elasticsearch 文档格式Elasticsearch Indexer: 将文档写入 Elasticsearch这种架构的优势在于低延迟、高可靠性且对数据库几乎无侵入。2. 增量索引更新实现增量索引更新是同步方案的核心主要通过 Canal 监听数据库变更事件并将变更应用到 Elasticsearch 索引中。实现增量更新的关键步骤2.1. 配置 Canal 监听数据库首先需要在 MySQL 数据库中开启 binlog 功能并配置 Canal 监听指定数据库。修改 my.cnf 文件添加以下配置[mysqld] server-id1 log-binmysql-bin binlog-formatROW binlog-row-imageFULL2.2. 创建 Canal 实例创建一个新的 Canal 实例指向需要同步的数据库canal.instance.mysql.slaveId1234 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal canal.instance.dbNameyour_database canal.instance.dbEncodingUTF-82.3. 实现消息处理逻辑编写 Canal Client 接收变更事件并将数据写入 Elasticsearchpublic class ElasticsearchHandler implements EntryHandlerCanalEntry.Entry { private RestHighLevelClient esClient; public ElasticsearchHandler(RestHighLevelClient esClient) { this.esClient esClient; } Override public void handle(CanalEntry.Entry entry) throws Exception { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (rowChange.getEventType() CanalEntry.EventType.INSERT || rowChange.getEventType() CanalEntry.EventType.UPDATE) { // 处理插入和更新操作 IndexRequest request new IndexRequest(your_index) .id(rowData.getAfterColumns(0).getValue()) .source(convertToMap(rowData.getAfterColumnsList())); esClient.index(request, RequestOptions.DEFAULT); } else if (rowChange.getEventType() CanalEntry.EventType.DELETE) { // 处理删除操作 DeleteRequest request new DeleteRequest(your_index) .id(rowData.getBeforeColumns(0).getValue()); esClient.delete(request, RequestOptions.DEFAULT); } } } } private MapString, Object convertToMap(ListCanalEntry.Column columns) { MapString, Object map new HashMap(); for (CanalEntry.Column column : columns) { if (column.getIsNull()) { map.put(column.getName(), null); } else { map.put(column.getName(), column.getValue()); } } return map; } }2.4. 处理批量提交为提高性能可以使用批量提交机制BulkRequest bulkRequest new BulkRequest(); // 添加多个索引/删除请求到批量请求中 // ... // 执行批量提交 BulkResponse bulkResponse esClient.bulk(bulkRequest, RequestOptions.DEFAULT); if (bulkResponse.hasFailures()) { // 处理失败情况 }3. 数据删除处理方案在 Canal 与 Elasticsearch 同步过程中处理删除操作是一个关键点。与数据更新不同删除操作需要特别注意数据一致性和同步延迟问题。3.1. 基于主键的删除处理最简单的删除方式是基于主键进行删除如上述代码所示。这种方法适用于每个表都有明确主键的情况。3.2. 软删除与硬删除根据业务需求可以选择软删除或硬删除硬删除直接从 Elasticsearch 中删除文档软删除在文档中标记为已删除而不是真正删除适合需要保留历史数据的场景3.3. 删除事件过滤在某些场景下可能需要过滤特定的删除事件Override public void handle(CanalEntry.Entry entry) throws Exception { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (rowChange.getEventType() CanalEntry.EventType.DELETE) { // 检查是否需要跳过此删除事件 if (shouldSkipDelete(rowData)) { continue; } // 执行删除操作 DeleteRequest request new DeleteRequest(your_index) .id(rowData.getBeforeColumns(0).getValue()); esClient.delete(request, RequestOptions.DEFAULT); } } } } private boolean shouldSkipDelete(CanalEntry.RowData rowData) { // 根据业务逻辑判断是否跳过删除 // 例如特定状态的数据不执行删除操作 for (CanalEntry.Column column : rowData.getBeforeColumnsList()) { if (status.equals(column.getName()) inactive.equals(column.getValue())) { return true; } } return false; }3.4. 删除操作的幂等性确保删除操作的幂等性非常重要特别是在网络不稳定或重试的情况下public void safeDelete(String index, String id) { try { // 检查文档是否存在 GetRequest getRequest new GetRequest(index, id); boolean exists esClient.exists(getRequest, RequestOptions.DEFAULT); if (exists) { // 文档存在则删除 DeleteRequest deleteRequest new DeleteRequest(index, id); esClient.delete(deleteRequest, RequestOptions.DEFAULT); } // 如果文档不存在不做任何操作 } catch (ElasticsearchException e) { // 处理异常如文档已被其他线程删除的情况 if (e.getDetailedMessage().contains(missing)) { // 文档不存在无需处理 return; } throw e; } }4. 全量重建策略虽然增量同步能够保持数据一致性但在某些情况下需要进行全量重建例如首次同步索引结构变更数据发生严重不一致需要修复4.1. 全量重建方案设计全量重建的基本流程如下停止 Canal 增量同步从数据库导出全量数据清空 Elasticsearch 索引将全量数据导入 Elasticsearch重新启动 Canal 增量同步4.2. 实现全量数据导出可以使用 JDBC 直接从数据库查询全量数据public ListMapString, Object exportFullData(String sql, Connection connection) throws SQLException { ListMapString, Object result new ArrayList(); try (PreparedStatement stmt connection.prepareStatement(sql); ResultSet rs stmt.executeQuery()) { ResultSetMetaData metaData rs.getMetaData(); int columnCount metaData.getColumnCount(); while (rs.next()) { MapString, Object row new LinkedHashMap(); for (int i 1; i columnCount; i) { row.put(metaData.getColumnName(i), rs.getObject(i)); } result.add(row); } } return result; }4.3. 批量导入 Elasticsearch使用 Elasticsearch 的批量 API 高效导入数据public void bulkIndexToES(ListMapString, Object documents, String index) throws IOException { BulkRequest bulkRequest new BulkRequest(); for (MapString, Object doc : documents) { // 假设文档中包含 id 字段 String id doc.get(id).toString(); // 移除 id 字段因为它在 IndexRequest 中单独指定 doc.remove(id); IndexRequest request new IndexRequest(index).id(id).source(doc); bulkRequest.add(request); // 每 1000 条提交一次 if (bulkRequest.numberOfActions() 1000) { BulkResponse bulkResponse esClient.bulk(bulkRequest, RequestOptions.DEFAULT); if (bulkResponse.hasFailures()) { // 处理失败 handleFailures(bulkResponse); } bulkRequest new BulkRequest(); } } // 提交剩余的请求 if (bulkRequest.numberOfActions() 0) { BulkResponse bulkResponse esClient.bulk(bulkRequest, RequestOptions.DEFAULT); if (bulkResponse.hasFailures()) { handleFailures(bulkResponse); } } } private void handleFailures(BulkResponse bulkResponse) { for (BulkItemResponse response : bulkResponse) { if (response.isFailed()) { // 记录失败信息 System.err.println(Failed to process document: response.getId() , Error: response.getFailure().getMessage()); } } }4.4. 全量重建与增量同步的衔接为避免数据丢失全量重建与增量同步之间需要正确衔接记录全量数据导出的时间点 T在时间点 T 之后的所有数据库变更需要单独记录全量数据导入完成后将时间点 T 之后的增量变更同步到 Elasticsearch这可以通过记录 binlog 位置来实现public class BinlogPosition { private String logFileName; private long logFileOffset; // 获取方法 public String getLogFileName() { return logFileName; } public long getLogFileOffset() { return logFileOffset; } // 设置方法 public void setLogFileName(String logFileName) { this.logFileName logFileName; } public void setLogFileOffset(long logFileOffset) { this.logFileOffset logFileOffset; } } // 在全量导出开始前记录位置 public BinlogPosition getCurrentBinlogPosition(Connection mysqlConn) throws SQLException { BinlogPosition position new BinlogPosition(); try (Statement stmt mysqlConn.createStatement(); ResultSet rs stmt.executeQuery(SHOW MASTER STATUS)) { if (rs.next()) { position.setLogFileName(rs.getString(File)); position.setLogFileOffset(rs.getLong(Position)); } } return position; } // 在全量导入完成后从此位置继续同步 public void resumeIncrementalSync(BinlogPosition position) { // 配置 Canal 从指定位置开始监听 canalConfig.setMasterId(position.getLogFileName()); canalConfig.setSlaveId(position.getLogFileOffset()); // 启动 Canal 客户端 startCanalClient(); }5. 实践案例与注意事项5.1. 完整的最小示例下面是一个完整的 Canal 与 Elasticsearch 同步的最小示例public class CanalElasticsearchSync { private RestHighLevelClient esClient; private CanalClient canalClient; public void init() { // 初始化 Elasticsearch 客户端 esClient new RestHighLevelClient( RestClient.builder(new HttpHost(localhost, 9200, http))); // 初始化 Canal 客户端 canalClient new CanalConnector(localhost, 11111, canal, canal, example); canalClient.connect(); canalClient.subscribe(); canalClient.rollback(); } public void startSync() { try { while (true) { Message message canalClient.getWithoutAck(100); long batchId message.getId(); if (message.getEntries() ! null !message.getEntries().isEmpty()) { for (CanalEntry.Entry entry : message.getEntries()) { handleEntry(entry); } } canalClient.ack(batchId); } } catch (Exception e) { e.printStackTrace(); } finally { canalClient.disconnect(); try { esClient.close(); } catch (IOException e) { e.printStackTrace(); } } } private void handleEntry(CanalEntry.Entry entry) throws Exception { if (entry.getEntryType() CanalEntry.EntryType.ROWDATA) { CanalEntry.RowChange rowChange CanalEntry.RowChange.parseFrom(entry.getStoreValue()); for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) { if (rowChange.getEventType() CanalEntry.EventType.INSERT || rowChange.getEventType() CanalEntry.EventType.UPDATE) { // 处理插入和更新 IndexRequest request new IndexRequest(your_index) .id(rowData.getAfterColumns(0).getValue()) .source(convertColumnsToMap(rowData.getAfterColumnsList())); esClient.index(request, RequestOptions.DEFAULT); } else if (rowChange.getEventType() CanalEntry.EventType.DELETE) { // 处理删除 DeleteRequest request new DeleteRequest(your_index) .id(rowData.getBeforeColumns(0).getValue()); esClient.delete(request, RequestOptions.DEFAULT); } } } } private MapString, Object convertColumnsToMap(ListCanalEntry.Column columns) { MapString, Object map new HashMap(); for (CanalEntry.Column column : columns) { if (column.getIsNull()) { map.put(column.getName(), null); } else { map.put(column.getName(), column.getValue()); } } return map; } public static void main(String[] args) { CanalElasticsearchSync sync new CanalElasticsearchSync(); sync.init(); sync.startSync(); } }5.2. 注意事项性能监控监控 Canal 和 Elasticsearch 的性能指标及时发现并处理性能瓶颈错误处理建立完善的错误处理机制特别是网络中断、数据格式错误等情况数据一致性定期检查 Canal 和 Elasticsearch 之间的数据一致性特别是关键业务数据备份策略制定并执行数据备份策略防止数据丢失版本兼容性确保 Canal 和 Elasticsearch 版本兼容避免因版本不匹配导致的问题资源管理合理配置内存和线程资源避免资源耗尽导致系统崩溃5.3. 同步策略对比| 同步策略 | 优点 | 缺点 | 适用场景 ||---------|------|------|---------|| 实时同步 | 低延迟数据最新 | 对数据库压力大资源消耗高 | 对实时性要求高的场景 || 批量同步 | 资源消耗小系统稳定性高 | 同步延迟高 | 对实时性要求不高的场景 || 混合策略 | 平衡实时性与资源消耗 | 实现复杂度高 | 大规模数据同步场景 |下面是 Canal 与 Elasticsearch 实时同步的流程图Binlog变更数据处理数据监控监控监控MySQL数据库Canal服务器消息队列/处理程序Elasticsearch监控工具
返回列表
PREV
查看更多资讯
NEXT
返回资讯列表