From 4df60a0fb58313b4fcb7e817849324eb8a83de0d Mon Sep 17 00:00:00 2001 From: mcy Date: Wed, 13 Mar 2019 16:57:23 +0800 Subject: [PATCH] =?UTF-8?q?es=E5=A2=9E=E5=8A=A0=E7=88=B6=E5=AD=90=E6=96=87?= =?UTF-8?q?=E6=A1=A3=E7=B4=A2=E5=BC=95=E5=90=8C=E6=AD=A5=E9=80=82=E9=85=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/es/config/ESSyncConfig.java | 47 +++----- .../adapter/es/service/ESSyncService.java | 6 +- .../client/adapter/es/support/ESTemplate.java | 109 ++++++++++++++++-- .../src/main/resources/es/biz_order.yml | 21 ++++ .../src/main/resources/es/customer.yml | 47 ++++++++ 5 files changed, 187 insertions(+), 43 deletions(-) create mode 100644 client-adapter/elasticsearch/src/main/resources/es/biz_order.yml create mode 100644 client-adapter/elasticsearch/src/main/resources/es/customer.yml diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/config/ESSyncConfig.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/config/ESSyncConfig.java index a84223fe..da37189b 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/config/ESSyncConfig.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/config/ESSyncConfig.java @@ -26,12 +26,6 @@ public class ESSyncConfig { private ESMapping esMapping; public void validate() { - if (!esMapping.relations.isEmpty()) { - Map relationMappings = esMapping.relations.get(0); - RelationMapping relationMapping = relationMappings.values().iterator().next(); - esMapping.isChild = StringUtils.isNotEmpty(relationMapping.getParent()); - } - if (esMapping._index == null) { throw new NullPointerException("esMapping._index"); } @@ -88,24 +82,23 @@ public class ESSyncConfig { public static class ESMapping { - private String _index; - private String _type; - private String _id; - private boolean upsert = false; - private String pk; - private List> relations = new ArrayList<>(); - private boolean isChild = false; - private String sql; + private String _index; + private String _type; + private String _id; + private boolean upsert = false; + private String pk; + private Map relations = new LinkedHashMap<>(); + private String sql; // 对象字段, 例: objFields: // - _labels: array:; - private Map objFields = new LinkedHashMap<>(); - private List skips = new ArrayList<>(); - private int commitBatch = 1000; - private String etlCondition; - private boolean syncByTimestamp = false; // 是否按时间戳定时同步 - private Long syncInterval; // 同步时间间隔 + private Map objFields = new LinkedHashMap<>(); + private List skips = new ArrayList<>(); + private int commitBatch = 1000; + private String etlCondition; + private boolean syncByTimestamp = false; // 是否按时间戳定时同步 + private Long syncInterval; // 同步时间间隔 - private SchemaItem schemaItem; // sql解析结果模型 + private SchemaItem schemaItem; // sql解析结果模型 public String get_index() { return _index; @@ -163,22 +156,14 @@ public class ESSyncConfig { this.skips = skips; } - public List> getRelations() { + public Map getRelations() { return relations; } - public void setRelations(List> relations) { + public void setRelations(Map relations) { this.relations = relations; } - public boolean isChild() { - return isChild; - } - - public void setChild(boolean child) { - isChild = child; - } - public String getSql() { return sql; } diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java index 0ba4d442..d5258223 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/service/ESSyncService.java @@ -43,7 +43,7 @@ public class ESSyncService { long begin = System.currentTimeMillis(); if (esSyncConfigs != null) { if (logger.isTraceEnabled()) { - logger.trace("Destination: {}, database:{}, table:{}, type:{}, effect index count: {}", + logger.trace("Destination: {}, database:{}, table:{}, type:{}, affected index count: {}", dml.getDestination(), dml.getDatabase(), dml.getTable(), @@ -65,7 +65,7 @@ public class ESSyncService { } } if (logger.isTraceEnabled()) { - logger.trace("Sync elapsed time: {} ms, effect index count:{}, destination: {}", + logger.trace("Sync elapsed time: {} ms, affected indexes count:{}, destination: {}", (System.currentTimeMillis() - begin), esSyncConfigs.size(), dml.getDestination()); @@ -74,7 +74,7 @@ public class ESSyncService { StringBuilder configIndexes = new StringBuilder(); esSyncConfigs .forEach(esSyncConfig -> configIndexes.append(esSyncConfig.getEsMapping().get_index()).append(" ")); - logger.debug("DML: {} \nEffect indexes: {}", + logger.debug("DML: {} \nAffected indexes: {}", JSON.toJSONString(dml, SerializerFeature.WriteMapNullValue), configIndexes.toString()); } diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESTemplate.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESTemplate.java index 997f1366..0a67ff78 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESTemplate.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESTemplate.java @@ -2,6 +2,7 @@ package com.alibaba.otter.canal.client.adapter.es.support; import java.sql.ResultSet; import java.sql.SQLException; +import java.util.HashMap; import java.util.LinkedHashMap; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -9,10 +10,13 @@ import java.util.concurrent.ConcurrentMap; import javax.sql.DataSource; +import org.apache.commons.lang.StringUtils; import org.elasticsearch.action.bulk.BulkItemResponse; import org.elasticsearch.action.bulk.BulkRequestBuilder; import org.elasticsearch.action.bulk.BulkResponse; +import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.action.search.SearchResponse; +import org.elasticsearch.action.update.UpdateRequestBuilder; import org.elasticsearch.client.transport.TransportClient; import org.elasticsearch.cluster.metadata.MappingMetaData; import org.elasticsearch.common.collect.ImmutableOpenMap; @@ -64,13 +68,24 @@ public class ESTemplate { */ public void insert(ESMapping mapping, Object pkVal, Map esFieldData) { if (mapping.get_id() != null) { + String parentVal = (String) esFieldData.remove("$parent_routing"); if (mapping.isUpsert()) { - getBulk().add(transportClient.prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) + UpdateRequestBuilder updateRequestBuilder = transportClient + .prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) .setDoc(esFieldData) - .setDocAsUpsert(true)); + .setDocAsUpsert(true); + if (StringUtils.isNotEmpty(parentVal)) { + updateRequestBuilder.setRouting(parentVal); + } + getBulk().add(updateRequestBuilder); } else { - getBulk().add(transportClient.prepareIndex(mapping.get_index(), mapping.get_type(), pkVal.toString()) - .setSource(esFieldData)); + IndexRequestBuilder indexRequestBuilder = transportClient + .prepareIndex(mapping.get_index(), mapping.get_type(), pkVal.toString()) + .setSource(esFieldData); + if (StringUtils.isNotEmpty(parentVal)) { + indexRequestBuilder.setRouting(parentVal); + } + getBulk().add(indexRequestBuilder); } commitBulk(); } else { @@ -137,7 +152,7 @@ public class ESTemplate { return count; }); if (logger.isTraceEnabled()) { - logger.trace("Update ES by query effect {} records", syncCount); + logger.trace("Update ES by query affected {} records", syncCount); } } @@ -200,13 +215,24 @@ public class ESTemplate { private void append4Update(ESMapping mapping, Object pkVal, Map esFieldData) { if (mapping.get_id() != null) { + String parentVal = (String) esFieldData.remove("$parent_routing"); if (mapping.isUpsert()) { - getBulk().add(transportClient.prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) + UpdateRequestBuilder updateRequestBuilder = transportClient + .prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) .setDoc(esFieldData) - .setDocAsUpsert(true)); + .setDocAsUpsert(true); + if (StringUtils.isNotEmpty(parentVal)) { + updateRequestBuilder.setRouting(parentVal); + } + getBulk().add(updateRequestBuilder); } else { - getBulk().add(transportClient.prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) - .setDoc(esFieldData)); + UpdateRequestBuilder updateRequestBuilder = transportClient + .prepareUpdate(mapping.get_index(), mapping.get_type(), pkVal.toString()) + .setDoc(esFieldData); + if (StringUtils.isNotEmpty(parentVal)) { + updateRequestBuilder.setRouting(parentVal); + } + getBulk().add(updateRequestBuilder); } } else { SearchResponse response = transportClient.prepareSearch(mapping.get_index()) @@ -257,6 +283,10 @@ public class ESTemplate { esFieldData.put(fieldItem.getFieldName(), value); } } + + // 添加父子文档关联信息 + putRelationDataFromRS(mapping, schemaItem, resultSet, esFieldData); + return resultIdVal; } @@ -294,6 +324,10 @@ public class ESTemplate { } } } + + // 添加父子文档关联信息 + putRelationDataFromRS(mapping, schemaItem, resultSet, esFieldData); + return resultIdVal; } @@ -340,6 +374,9 @@ public class ESTemplate { esFieldData.put(fieldItem.getFieldName(), value); } } + + // 添加父子文档关联信息 + putRelationData(mapping, schemaItem, dmlData, esFieldData); return resultIdVal; } @@ -368,9 +405,63 @@ public class ESTemplate { getValFromData(mapping, dmlData, fieldItem.getFieldName(), columnName)); } } + + // 添加父子文档关联信息 + putRelationData(mapping, schemaItem, dmlOld, esFieldData); return resultIdVal; } + private void putRelationDataFromRS(ESMapping mapping, SchemaItem schemaItem, ResultSet resultSet, + Map esFieldData) { + // 添加父子文档关联信息 + if (!mapping.getRelations().isEmpty()) { + mapping.getRelations().forEach((relationField, relationMapping) -> { + Map relations = new HashMap<>(); + relations.put("name", relationMapping.getName()); + if (StringUtils.isNotEmpty(relationMapping.getParent())) { + FieldItem parentFieldItem = schemaItem.getSelectFields().get(relationMapping.getParent()); + Object parentVal; + try { + parentVal = getValFromRS(mapping, + resultSet, + parentFieldItem.getFieldName(), + parentFieldItem.getFieldName()); + } catch (SQLException e) { + throw new RuntimeException(e); + } + if (parentVal != null) { + relations.put("parent", parentVal.toString()); + esFieldData.put("$parent_routing", parentVal.toString()); + + } + } + esFieldData.put(relationField, relations); + }); + } + } + + private void putRelationData(ESMapping mapping, SchemaItem schemaItem, Map dmlData, + Map esFieldData) { + // 添加父子文档关联信息 + if (!mapping.getRelations().isEmpty()) { + mapping.getRelations().forEach((relationField, relationMapping) -> { + Map relations = new HashMap<>(); + relations.put("name", relationMapping.getName()); + if (StringUtils.isNotEmpty(relationMapping.getParent())) { + FieldItem parentFieldItem = schemaItem.getSelectFields().get(relationMapping.getParent()); + String columnName = parentFieldItem.getColumnItems().iterator().next().getColumnName(); + Object parentVal = getValFromData(mapping, dmlData, parentFieldItem.getFieldName(), columnName); + if (parentVal != null) { + relations.put("parent", parentVal.toString()); + esFieldData.put("$parent_routing", parentVal.toString()); + + } + } + esFieldData.put(relationField, relations); + }); + } + } + /** * es 字段类型本地缓存 */ diff --git a/client-adapter/elasticsearch/src/main/resources/es/biz_order.yml b/client-adapter/elasticsearch/src/main/resources/es/biz_order.yml new file mode 100644 index 00000000..c49806fa --- /dev/null +++ b/client-adapter/elasticsearch/src/main/resources/es/biz_order.yml @@ -0,0 +1,21 @@ +dataSourceKey: defaultDS +destination: example +groupId: g1 +esMapping: + _index: customer + _type: _doc + _id: _id + relations: + customer_order: + name: order + parent: customer_id + sql: "select concat('oid_', t.id) as _id, + t.customer_id, + t.id as order_id, + t.serial_code as order_serial, + t.c_time as order_time + from biz_order t" + skips: + - customer_id + etlCondition: "where t.c_time>='{0}'" + commitBatch: 3000 diff --git a/client-adapter/elasticsearch/src/main/resources/es/customer.yml b/client-adapter/elasticsearch/src/main/resources/es/customer.yml new file mode 100644 index 00000000..06793a58 --- /dev/null +++ b/client-adapter/elasticsearch/src/main/resources/es/customer.yml @@ -0,0 +1,47 @@ +dataSourceKey: defaultDS +destination: example +groupId: g1 +esMapping: + _index: customer + _type: _doc + _id: id + relations: + customer_order: + name: customer + sql: "select t.id, t.name, t.email from customer t" + etlCondition: "where t.c_time>='{0}'" + commitBatch: 3000 + + +#{ +# "mappings":{ +# "_doc":{ +# "properties":{ +# "id": { +# "type": "long" +# }, +# "name": { +# "type": "text" +# }, +# "email": { +# "type": "text" +# }, +# "order_id": { +# "type": "long" +# }, +# "order_serial": { +# "type": "text" +# }, +# "order_time": { +# "type": "date" +# }, +# "customer_order":{ +# "type":"join", +# "relations":{ +# "customer":"order" +# } +# } +# } +# } +# } +#}