独行的浪子
2022-04-26 3ad94b9abcbde117c6d6ce43a9d023b7755ec0eb
提交 | 用户 | age
de8c2b 1 package com.duxinglangzi.canal.starter.container;
D 2
3 import com.alibaba.otter.canal.client.CanalConnector;
4 import com.alibaba.otter.canal.protocol.CanalEntry;
5 import com.alibaba.otter.canal.protocol.Message;
6 import com.duxinglangzi.canal.starter.configuration.CanalAutoConfigurationProperties;
7 import com.duxinglangzi.canal.starter.configuration.CanalListenerEndpointRegistrar;
8 import org.slf4j.Logger;
9 import org.slf4j.LoggerFactory;
10
11 import java.lang.reflect.InvocationTargetException;
12 import java.util.*;
13
14 /**
627420 15  * DML 数据拉取、解析
de8c2b 16  * @author wuqiong 2022/4/11
D 17  */
18 public class DmlMessageTransponderContainer extends AbstractCanalTransponderContainer {
19
20     private final static Logger logger = LoggerFactory.getLogger(DmlMessageTransponderContainer.class);
21
22     private CanalConnector connector;
23     private CanalAutoConfigurationProperties.EndpointInstance endpointInstance;
24     private List<CanalListenerEndpointRegistrar> registrars = new ArrayList<>();
25     private Set<CanalEntry.EventType> SUPPORT_ALL_TYPES = new HashSet<>();
26
27
28     public void initConnect() {
29         // init supportAllTypes
30         registrars.forEach(e -> SUPPORT_ALL_TYPES.addAll(
31                 Arrays.asList(e.getListenerEntry().getValue().eventType())));
32         connector.connect();
33         connector.subscribe(endpointInstance.getSubscribe());
34         connector.rollback();
35
36     }
37
929202 38     public void disconnect(){
D 39         // 关闭连接
40         connector.disconnect();
41     }
42
de8c2b 43
D 44     public void doStart() {
45         Message message = null;
46         try {
47             message = connector.getWithoutAck(endpointInstance.getBatchSize()); // 获取指定数量的数据
48         } catch (Exception clientException) {
49             logger.error("[MessageTransponderContainer] error msg : ", clientException);
50             endpointInstance.setRetryCount(endpointInstance.getRetryCount() - 1);
51             if (endpointInstance.getRetryCount() < 0) {
52                 logger.error("[MessageTransponderContainer] retry count is zero ,  " +
53                                 "thread interrupt , current connector host: {} , port: {} ",
54                         endpointInstance.getHost(), endpointInstance.getPort());
55                 Thread.currentThread().interrupt();
929202 56                 disconnect();
de8c2b 57             } else {
D 58                 sleep(endpointInstance.getAcquireInterval());
59             }
60             return;
61         }
62         List<CanalEntry.Entry> entries = message.getEntries();
63         if (message.getId() == -1 || entries.isEmpty()) {
64             sleep(endpointInstance.getAcquireInterval());
65             return;
66         }
3ad94b 67         for (CanalEntry.Entry entry : entries) {
68             try {
69                 consumer(entry);
70             } catch (Exception e) {
71                 e.printStackTrace();
72                 logger.error("[MessageTransponderContainer_doStart] CanalEntry.Entry consumer error ", e);
73                 // connector.rollback(message.getId()); // 目前先不处理失败, 无需回滚数据
74                 // return;
75             }
de8c2b 76         }
D 77         connector.ack(message.getId()); // 提交确认
78     }
79
80     public DmlMessageTransponderContainer(
81             CanalConnector connector, List<CanalListenerEndpointRegistrar> registrars,
82             CanalAutoConfigurationProperties.EndpointInstance endpointInstance) {
83         this.connector = connector;
84         this.registrars.addAll(registrars);
85         this.endpointInstance = endpointInstance;
86     }
87
88     private void consumer(CanalEntry.Entry entry) {
89         if (IGNORE_ENTRY_TYPES.contains(entry.getEntryType())) return;
90         CanalEntry.RowChange rowChange = null;
91         try {
92             rowChange = CanalEntry.RowChange.parseFrom(entry.getStoreValue());
93         } catch (Exception e) {
94             logger.error("[MessageTransponderContainer_consumer] RowChange parse has an error ", e);
95             throw new RuntimeException("RowChange parse has an error , data:" + entry.toString(), e);
96         }
97
98         // 忽略 ddl 语句
99         if (rowChange.hasIsDdl() && rowChange.getIsDdl()) return;
100         CanalEntry.EventType eventType = rowChange.getEventType();
101         if (!SUPPORT_ALL_TYPES.contains(eventType)) return;
102         for (CanalEntry.RowData rowData : rowChange.getRowDatasList()) {
103             registrars
104                     .stream()
105                     .filter(CanalListenerEndpointRegistrar.filterArgs(
106                             entry.getHeader().getSchemaName(),
107                             entry.getHeader().getTableName(),
108                             eventType))
109                     .forEach(element -> {
110                         try {
111                             element.getListenerEntry().getKey().invoke(element.getBean(), eventType, rowData);
112                         } catch (IllegalAccessException | InvocationTargetException e) {
113                             logger.error("[MessageTransponderContainer_consumer] RowData Callback Method invoke error message", e);
114                             throw new RuntimeException("RowData Callback Method invoke error message: " + e.getMessage(), e);
115                         }
116                     });
117         }
118     }
119
120 }