時(shí)同步:增量索引更新、刪除處理與全量重建方案)
1. Canal 與 Elasticsearch 實(shí)時(shí)同步概述Canal 是阿里巴巴開(kāi)源的基于數(shù)據(jù)庫(kù)增量日志解析的組件支持 MySQL、Oracle 等數(shù)據(jù)庫(kù)。它通過(guò)解析數(shù)據(jù)庫(kù)的 binlog 日志將數(shù)據(jù)變更事件推送到消息隊(duì)列或直接處理實(shí)現(xiàn)數(shù)據(jù)的準(zhǔn)實(shí)時(shí)同步。Elasticsearch 是一個(gè)基于 Lucene 庫(kù)的搜索引擎具有強(qiáng)大的全文檢索和分析能力。將 Canal 與 Elasticsearch 結(jié)合使用可以實(shí)現(xiàn)數(shù)據(jù)庫(kù)數(shù)據(jù)到 Elasticsearch 的準(zhǔn)實(shí)時(shí)同步充分發(fā)揮 Elasticsearch 在搜索、分析和日志處理方面的優(yōu)勢(shì)。這種架構(gòu)廣泛應(yīng)用于電商搜索、日志分析、監(jiān)控告警等場(chǎng)景。Canal 與 Elasticsearch 同步的基本架構(gòu)包括Canal Server: 監(jiān)聽(tīng)并解析數(shù)據(jù)庫(kù) binlogCanal Client: 接收變更事件并處理數(shù)據(jù)轉(zhuǎn)換: 將數(shù)據(jù)庫(kù)數(shù)據(jù)轉(zhuǎn)換為 Elasticsearch 文檔格式Elasticsearch Indexer: 將文檔寫(xiě)入 Elasticsearch這種架構(gòu)的優(yōu)勢(shì)在于低延遲、高可靠性且對(duì)數(shù)據(jù)庫(kù)幾乎無(wú)侵入。2. 增量索引更新實(shí)現(xiàn)增量索引更新是同步方案的核心主要通過(guò) Canal 監(jiān)聽(tīng)數(shù)據(jù)庫(kù)變更事件并將變更應(yīng)用到 Elasticsearch 索引中。實(shí)現(xiàn)增量更新的關(guān)鍵步驟2.1. 配置 Canal 監(jiān)聽(tīng)數(shù)據(jù)庫(kù)首先需要在 MySQL 數(shù)據(jù)庫(kù)中開(kāi)啟 binlog 功能并配置 Canal 監(jiān)聽(tīng)指定數(shù)據(jù)庫(kù)。修改 my.cnf 文件添加以下配置[mysqld] server-id1 log-binmysql-bin binlog-formatROW binlog-row-imageFULL2.2. 創(chuàng)建 Canal 實(shí)例創(chuàng)建一個(gè)新的 Canal 實(shí)例指向需要同步的數(shù)據(jù)庫(kù)canal.instance.mysql.slaveId1234 canal.instance.dbUsernamecanal canal.instance.dbPasswordcanal canal.instance.dbNameyour_database canal.instance.dbEncodingUTF-82.3. 實(shí)現(xiàn)消息處理邏輯編寫(xiě) Canal Client 接收變更事件并將數(shù)據(jù)寫(xiě)入 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. 處理批量提交為提高性能可以使用批量提交機(jī)制BulkRequest bulkRequest new BulkRequest(); // 添加多個(gè)索引/刪除請(qǐng)求到批量請(qǐng)求中 // ... // 執(zhí)行批量提交 BulkResponse bulkResponse esClient.bulk(bulkRequest, RequestOptions.DEFAULT); if (bulkResponse.hasFailures()) { // 處理失敗情況 }3. 數(shù)據(jù)刪除處理方案在 Canal 與 Elasticsearch 同步過(guò)程中處理刪除操作是一個(gè)關(guān)鍵點(diǎn)。與數(shù)據(jù)更新不同刪除操作需要特別注意數(shù)據(jù)一致性和同步延遲問(wèn)題。3.1. 基于主鍵的刪除處理最簡(jiǎn)單的刪除方式是基于主鍵進(jìn)行刪除如上述代碼所示。這種方法適用于每個(gè)表都有明確主鍵的情況。3.2. 軟刪除與硬刪除根據(jù)業(yè)務(wù)需求可以選擇軟刪除或硬刪除硬刪除直接從 Elasticsearch 中刪除文檔軟刪除在文檔中標(biāo)記為已刪除而不是真正刪除適合需要保留歷史數(shù)據(jù)的場(chǎng)景3.3. 刪除事件過(guò)濾在某些場(chǎng)景下可能需要過(guò)濾特定的刪除事件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) { // 檢查是否需要跳過(guò)此刪除事件 if (shouldSkipDelete(rowData)) { continue; } // 執(zhí)行刪除操作 DeleteRequest request new DeleteRequest(your_index) .id(rowData.getBeforeColumns(0).getValue()); esClient.delete(request, RequestOptions.DEFAULT); } } } } private boolean shouldSkipDelete(CanalEntry.RowData rowData) { // 根據(jù)業(yè)務(wù)邏輯判斷是否跳過(guò)刪除 // 例如特定狀態(tài)的數(shù)據(jù)不執(zhí)行刪除操作 for (CanalEntry.Column column : rowData.getBeforeColumnsList()) { if (status.equals(column.getName()) inactive.equals(column.getValue())) { return true; } } return false; }3.4. 刪除操作的冪等性確保刪除操作的冪等性非常重要特別是在網(wǎng)絡(luò)不穩(wěn)定或重試的情況下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)) { // 文檔不存在無(wú)需處理 return; } throw e; } }4. 全量重建策略雖然增量同步能夠保持?jǐn)?shù)據(jù)一致性但在某些情況下需要進(jìn)行全量重建例如首次同步索引結(jié)構(gòu)變更數(shù)據(jù)發(fā)生嚴(yán)重不一致需要修復(fù)4.1. 全量重建方案設(shè)計(jì)全量重建的基本流程如下停止 Canal 增量同步從數(shù)據(jù)庫(kù)導(dǎo)出全量數(shù)據(jù)清空 Elasticsearch 索引將全量數(shù)據(jù)導(dǎo)入 Elasticsearch重新啟動(dòng) Canal 增量同步4.2. 實(shí)現(xiàn)全量數(shù)據(jù)導(dǎo)出可以使用 JDBC 直接從數(shù)據(jù)庫(kù)查詢?nèi)繑?shù)據(jù)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. 批量導(dǎo)入 Elasticsearch使用 Elasticsearch 的批量 API 高效導(dǎo)入數(shù)據(jù)public void bulkIndexToES(ListMapString, Object documents, String index) throws IOException { BulkRequest bulkRequest new BulkRequest(); for (MapString, Object doc : documents) { // 假設(shè)文檔中包含 id 字段 String id doc.get(id).toString(); // 移除 id 字段因?yàn)樗?IndexRequest 中單獨(dú)指定 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(); } } // 提交剩余的請(qǐng)求 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. 全量重建與增量同步的銜接為避免數(shù)據(jù)丟失全量重建與增量同步之間需要正確銜接記錄全量數(shù)據(jù)導(dǎo)出的時(shí)間點(diǎn) T在時(shí)間點(diǎn) T 之后的所有數(shù)據(jù)庫(kù)變更需要單獨(dú)記錄全量數(shù)據(jù)導(dǎo)入完成后將時(shí)間點(diǎn) T 之后的增量變更同步到 Elasticsearch這可以通過(guò)記錄 binlog 位置來(lái)實(shí)現(xiàn)public class BinlogPosition { private String logFileName; private long logFileOffset; // 獲取方法 public String getLogFileName() { return logFileName; } public long getLogFileOffset() { return logFileOffset; } // 設(shè)置方法 public void setLogFileName(String logFileName) { this.logFileName logFileName; } public void setLogFileOffset(long logFileOffset) { this.logFileOffset logFileOffset; } } // 在全量導(dǎo)出開(kāi)始前記錄位置 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; } // 在全量導(dǎo)入完成后從此位置繼續(xù)同步 public void resumeIncrementalSync(BinlogPosition position) { // 配置 Canal 從指定位置開(kāi)始監(jiān)聽(tīng) canalConfig.setMasterId(position.getLogFileName()); canalConfig.setSlaveId(position.getLogFileOffset()); // 啟動(dòng) Canal 客戶端 startCanalClient(); }5. 實(shí)踐案例與注意事項(xiàng)5.1. 完整的最小示例下面是一個(gè)完整的 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. 注意事項(xiàng)性能監(jiān)控監(jiān)控 Canal 和 Elasticsearch 的性能指標(biāo)及時(shí)發(fā)現(xiàn)并處理性能瓶頸錯(cuò)誤處理建立完善的錯(cuò)誤處理機(jī)制特別是網(wǎng)絡(luò)中斷、數(shù)據(jù)格式錯(cuò)誤等情況數(shù)據(jù)一致性定期檢查 Canal 和 Elasticsearch 之間的數(shù)據(jù)一致性特別是關(guān)鍵業(yè)務(wù)數(shù)據(jù)備份策略制定并執(zhí)行數(shù)據(jù)備份策略防止數(shù)據(jù)丟失版本兼容性確保 Canal 和 Elasticsearch 版本兼容避免因版本不匹配導(dǎo)致的問(wèn)題資源管理合理配置內(nèi)存和線程資源避免資源耗盡導(dǎo)致系統(tǒng)崩潰5.3. 同步策略對(duì)比| 同步策略 | 優(yōu)點(diǎn) | 缺點(diǎn) | 適用場(chǎng)景 ||---------|------|------|---------|| 實(shí)時(shí)同步 | 低延遲數(shù)據(jù)最新 | 對(duì)數(shù)據(jù)庫(kù)壓力大資源消耗高 | 對(duì)實(shí)時(shí)性要求高的場(chǎng)景 || 批量同步 | 資源消耗小系統(tǒng)穩(wěn)定性高 | 同步延遲高 | 對(duì)實(shí)時(shí)性要求不高的場(chǎng)景 || 混合策略 | 平衡實(shí)時(shí)性與資源消耗 | 實(shí)現(xiàn)復(fù)雜度高 | 大規(guī)模數(shù)據(jù)同步場(chǎng)景 |下面是 Canal 與 Elasticsearch 實(shí)時(shí)同步的流程圖Binlog變更數(shù)據(jù)處理數(shù)據(jù)監(jiān)控監(jiān)控監(jiān)控MySQL數(shù)據(jù)庫(kù)Canal服務(wù)器消息隊(duì)列/處理程序Elasticsearch監(jiān)控工具