IncrementPuller.java 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328
  1. package org.dbsyncer.manager.puller;
  2. import org.dbsyncer.common.event.ChangedEvent;
  3. import org.dbsyncer.common.event.ChangedOffset;
  4. import org.dbsyncer.common.event.CommonChangedEvent;
  5. import org.dbsyncer.common.event.DDLChangedEvent;
  6. import org.dbsyncer.common.event.RefreshOffsetEvent;
  7. import org.dbsyncer.common.event.RowChangedEvent;
  8. import org.dbsyncer.common.event.ScanChangedEvent;
  9. import org.dbsyncer.common.event.Watcher;
  10. import org.dbsyncer.common.model.AbstractConnectorConfig;
  11. import org.dbsyncer.common.scheduled.ScheduledTaskJob;
  12. import org.dbsyncer.common.scheduled.ScheduledTaskService;
  13. import org.dbsyncer.common.util.CollectionUtils;
  14. import org.dbsyncer.connector.ConnectorFactory;
  15. import org.dbsyncer.connector.model.Table;
  16. import org.dbsyncer.listener.AbstractExtractor;
  17. import org.dbsyncer.listener.Extractor;
  18. import org.dbsyncer.listener.Listener;
  19. import org.dbsyncer.listener.config.ListenerConfig;
  20. import org.dbsyncer.listener.enums.ListenerTypeEnum;
  21. import org.dbsyncer.listener.quartz.AbstractQuartzExtractor;
  22. import org.dbsyncer.listener.quartz.TableGroupQuartzCommand;
  23. import org.dbsyncer.manager.Manager;
  24. import org.dbsyncer.manager.ManagerException;
  25. import org.dbsyncer.manager.model.FieldPicker;
  26. import org.dbsyncer.parser.logger.LogService;
  27. import org.dbsyncer.parser.logger.LogType;
  28. import org.dbsyncer.parser.model.Connector;
  29. import org.dbsyncer.parser.model.Mapping;
  30. import org.dbsyncer.parser.model.Meta;
  31. import org.dbsyncer.parser.model.TableGroup;
  32. import org.dbsyncer.parser.util.PickerUtil;
  33. import org.slf4j.Logger;
  34. import org.slf4j.LoggerFactory;
  35. import org.springframework.context.ApplicationListener;
  36. import org.springframework.stereotype.Component;
  37. import org.springframework.util.Assert;
  38. import javax.annotation.PostConstruct;
  39. import javax.annotation.Resource;
  40. import java.time.Instant;
  41. import java.util.ArrayList;
  42. import java.util.HashSet;
  43. import java.util.LinkedHashMap;
  44. import java.util.LinkedList;
  45. import java.util.List;
  46. import java.util.Map;
  47. import java.util.Set;
  48. import java.util.concurrent.ConcurrentHashMap;
  49. import java.util.function.Consumer;
  50. import java.util.stream.Collectors;
  51. /**
  52. * 增量同步
  53. *
  54. * @author AE86
  55. * @version 1.0.0
  56. * @date 2020/04/26 15:28
  57. */
  58. @Component
  59. public final class IncrementPuller extends AbstractPuller implements ApplicationListener<RefreshOffsetEvent>, ScheduledTaskJob {
  60. private final Logger logger = LoggerFactory.getLogger(getClass());
  61. @Resource
  62. private Listener listener;
  63. @Resource
  64. private Manager manager;
  65. @Resource
  66. private LogService logService;
  67. @Resource
  68. private ScheduledTaskService scheduledTaskService;
  69. @Resource
  70. private ConnectorFactory connectorFactory;
  71. @Resource
  72. private BufferActuatorRouter bufferActuatorRouter;
  73. private Map<String, Extractor> map = new ConcurrentHashMap<>();
  74. @PostConstruct
  75. private void init() {
  76. scheduledTaskService.start(3000, this);
  77. }
  78. @Override
  79. public void start(Mapping mapping) {
  80. final String mappingId = mapping.getId();
  81. final String metaId = mapping.getMetaId();
  82. logger.info("开始增量同步:{}, {}", metaId, mapping.getName());
  83. Connector connector = manager.getConnector(mapping.getSourceConnectorId());
  84. Assert.notNull(connector, "连接器不能为空.");
  85. List<TableGroup> list = manager.getSortedTableGroupAll(mappingId);
  86. Assert.notEmpty(list, "映射关系不能为空.");
  87. Meta meta = manager.getMeta(metaId);
  88. Assert.notNull(meta, "Meta不能为空.");
  89. Thread worker = new Thread(() -> {
  90. try {
  91. long now = Instant.now().toEpochMilli();
  92. meta.setBeginTime(now);
  93. meta.setEndTime(now);
  94. manager.editConfigModel(meta);
  95. map.putIfAbsent(metaId, getExtractor(mapping, connector, list, meta));
  96. map.get(metaId).start();
  97. } catch (Exception e) {
  98. close(metaId);
  99. logService.log(LogType.TableGroupLog.INCREMENT_FAILED, e.getMessage());
  100. logger.error("运行异常,结束增量同步{}:{}", metaId, e.getMessage());
  101. }
  102. });
  103. worker.setName(new StringBuilder("increment-worker-").append(mapping.getId()).toString());
  104. worker.setDaemon(false);
  105. worker.start();
  106. }
  107. @Override
  108. public void close(String metaId) {
  109. Extractor extractor = map.get(metaId);
  110. if (null != extractor) {
  111. bufferActuatorRouter.unbind(metaId);
  112. extractor.close();
  113. }
  114. map.remove(metaId);
  115. publishClosedEvent(metaId);
  116. logger.info("关闭成功:{}", metaId);
  117. }
  118. @Override
  119. public void onApplicationEvent(RefreshOffsetEvent event) {
  120. List<ChangedOffset> offsetList = event.getOffsetList();
  121. if (!CollectionUtils.isEmpty(offsetList)) {
  122. offsetList.forEach(offset -> {
  123. if (offset.isRefreshOffset() && map.containsKey(offset.getMetaId())) {
  124. map.get(offset.getMetaId()).refreshEvent(offset);
  125. }
  126. });
  127. }
  128. }
  129. @Override
  130. public void run() {
  131. // 定时同步增量信息
  132. map.values().forEach(extractor -> extractor.flushEvent());
  133. }
  134. private AbstractExtractor getExtractor(Mapping mapping, Connector connector, List<TableGroup> list, Meta meta) throws InstantiationException, IllegalAccessException {
  135. AbstractConnectorConfig connectorConfig = connector.getConfig();
  136. ListenerConfig listenerConfig = mapping.getListener();
  137. // timing/log
  138. final String listenerType = listenerConfig.getListenerType();
  139. AbstractExtractor extractor = null;
  140. // 默认定时抽取
  141. if (ListenerTypeEnum.isTiming(listenerType)) {
  142. AbstractQuartzExtractor quartzExtractor = listener.getExtractor(ListenerTypeEnum.TIMING, connectorConfig.getConnectorType(), AbstractQuartzExtractor.class);
  143. quartzExtractor.setCommands(list.stream().map(t -> new TableGroupQuartzCommand(t.getSourceTable(), t.getCommand())).collect(Collectors.toList()));
  144. quartzExtractor.register(new QuartzConsumer(meta, mapping, list));
  145. extractor = quartzExtractor;
  146. }
  147. // 基于日志抽取
  148. if (ListenerTypeEnum.isLog(listenerType)) {
  149. extractor = listener.getExtractor(ListenerTypeEnum.LOG, connectorConfig.getConnectorType(), AbstractExtractor.class);
  150. extractor.register(new LogConsumer(meta, mapping, list));
  151. }
  152. if (null != extractor) {
  153. Set<String> filterTable = new HashSet<>();
  154. List<Table> sourceTable = new ArrayList<>();
  155. list.forEach(t -> {
  156. Table table = t.getSourceTable();
  157. if (!filterTable.contains(t.getName())) {
  158. sourceTable.add(table);
  159. }
  160. filterTable.add(table.getName());
  161. });
  162. extractor.setConnectorFactory(connectorFactory);
  163. extractor.setScheduledTaskService(scheduledTaskService);
  164. extractor.setConnectorConfig(connectorConfig);
  165. extractor.setListenerConfig(listenerConfig);
  166. extractor.setFilterTable(filterTable);
  167. extractor.setSourceTable(sourceTable);
  168. extractor.setSnapshot(meta.getSnapshot());
  169. extractor.setMetaId(meta.getId());
  170. return extractor;
  171. }
  172. throw new ManagerException("未知的监听配置.");
  173. }
  174. abstract class AbstractConsumer<E extends ChangedEvent> implements Watcher {
  175. protected Meta meta;
  176. public abstract void onChange(E e);
  177. public void onDDLChanged(DDLChangedEvent event) {
  178. }
  179. @Override
  180. public void changeEvent(ChangedEvent event) {
  181. event.getChangedOffset().setMetaId(meta.getId());
  182. if (event instanceof DDLChangedEvent) {
  183. onDDLChanged((DDLChangedEvent) event);
  184. return;
  185. }
  186. onChange((E) event);
  187. }
  188. @Override
  189. public void flushEvent(Map<String, String> snapshot) {
  190. meta.setSnapshot(snapshot);
  191. manager.editConfigModel(meta);
  192. }
  193. @Override
  194. public void errorEvent(Exception e) {
  195. logService.log(LogType.TableGroupLog.INCREMENT_FAILED, e.getMessage());
  196. }
  197. @Override
  198. public long getMetaUpdateTime() {
  199. return meta.getUpdateTime();
  200. }
  201. protected void bind(String tableGroupId) {
  202. bufferActuatorRouter.bind(meta.getId(), tableGroupId);
  203. }
  204. protected void execute(String tableGroupId, ChangedEvent event) {
  205. bufferActuatorRouter.execute(meta.getId(), tableGroupId, event);
  206. }
  207. }
  208. final class QuartzConsumer extends AbstractConsumer<ScanChangedEvent> {
  209. private List<FieldPicker> tablePicker = new LinkedList<>();
  210. public QuartzConsumer(Meta meta, Mapping mapping, List<TableGroup> tableGroups) {
  211. this.meta = meta;
  212. tableGroups.forEach(t -> {
  213. tablePicker.add(new FieldPicker(PickerUtil.mergeTableGroupConfig(mapping, t)));
  214. bind(t.getId());
  215. });
  216. }
  217. @Override
  218. public void onChange(ScanChangedEvent event) {
  219. final FieldPicker picker = tablePicker.get(event.getTableGroupIndex());
  220. TableGroup tableGroup = picker.getTableGroup();
  221. event.setSourceTableName(tableGroup.getSourceTable().getName());
  222. // 定时暂不支持触发刷新增量点事件
  223. execute(tableGroup.getId(), event);
  224. }
  225. }
  226. final class LogConsumer extends AbstractConsumer<RowChangedEvent> {
  227. private Mapping mapping;
  228. private List<TableGroup> tableGroups;
  229. private Map<String, List<FieldPicker>> tablePicker = new LinkedHashMap<>();
  230. //判断上次是否为ddl,是ddl需要强制刷新下picker
  231. private boolean ddlChanged;
  232. public LogConsumer(Meta meta, Mapping mapping, List<TableGroup> tableGroups) {
  233. this.tableGroups = tableGroups;
  234. this.meta = meta;
  235. this.mapping = mapping;
  236. addTablePicker(true);
  237. }
  238. @Override
  239. public void onChange(RowChangedEvent event) {
  240. // 需要强制刷新 fix https://gitee.com/ghi/dbsyncer/issues/I8DJUR
  241. if (ddlChanged) {
  242. addTablePicker(false);
  243. ddlChanged = false;
  244. }
  245. process(event, picker -> {
  246. final Map<String, Object> changedRow = picker.getColumns(event.getDataList());
  247. if (picker.filter(changedRow)) {
  248. event.setChangedRow(changedRow);
  249. execute(picker.getTableGroup().getId(), event);
  250. }
  251. });
  252. }
  253. @Override
  254. public void onDDLChanged(DDLChangedEvent event) {
  255. ddlChanged = true;
  256. process(event, picker -> execute(picker.getTableGroup().getId(), event));
  257. }
  258. private void process(CommonChangedEvent event, Consumer<FieldPicker> consumer) {
  259. // 处理过程有异常向上抛
  260. List<FieldPicker> pickers = tablePicker.get(event.getSourceTableName());
  261. if (!CollectionUtils.isEmpty(pickers)) {
  262. // 触发刷新增量点事件
  263. event.getChangedOffset().setRefreshOffset(true);
  264. pickers.forEach(picker -> consumer.accept(picker));
  265. }
  266. }
  267. private void addTablePicker(boolean bindBufferActuatorRouter) {
  268. this.tablePicker.clear();
  269. this.tableGroups.forEach(t -> {
  270. final Table table = t.getSourceTable();
  271. final String tableName = table.getName();
  272. tablePicker.putIfAbsent(tableName, new ArrayList<>());
  273. TableGroup group = PickerUtil.mergeTableGroupConfig(mapping, t);
  274. tablePicker.get(tableName).add(new FieldPicker(group, group.getFilter(), table.getColumn(), group.getFieldMapping()));
  275. // 是否注册到路由服务中
  276. if (bindBufferActuatorRouter) {
  277. bind(group.getId());
  278. }
  279. });
  280. }
  281. }
  282. }