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 3c50d5c6..c2495416 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 @@ -102,6 +102,82 @@ public class ESSyncService { } + /** + * 插入操作dml + * + * @param config es配置 + * @param dml dml数据 + */ + private void insert(ESSyncConfig config, Dml dml) { + List> dataList = dml.getData(); + if (dataList == null || dataList.isEmpty()) { + return; + } + SchemaItem schemaItem = config.getEsMapping().getSchemaItem(); + for (Map data : dataList) { + if (data == null || data.isEmpty()) { + continue; + } + + if (schemaItem.getAliasTableItems().size() == 1 && schemaItem.isAllFieldsSimple()) { + // ------单表 & 所有字段都为简单字段------ + singleTableSimpleFiledInsert(config, dml, data); + } else { + // ------是主表 查询sql来插入------ + if (schemaItem.getMainTable().getTableName().equalsIgnoreCase(dml.getTable())) { + mainTableInsert(config, dml, data); + } + + // 从表的操作 + for (TableItem tableItem : schemaItem.getAliasTableItems().values()) { + if (tableItem.isMain()) { + continue; + } + if (!tableItem.getTableName().equals(dml.getTable())) { + continue; + } + // 关联条件出现在主表查询条件是否为简单字段 + boolean allFieldsSimple = true; + for (FieldItem fieldItem : tableItem.getRelationSelectFieldItems()) { + if (fieldItem.isMethod() || fieldItem.isBinaryOp()) { + allFieldsSimple = false; + break; + } + } + // 所有查询字段均为简单字段 + if (allFieldsSimple) { + // 不是子查询 + if (!tableItem.isSubQuery()) { + // ------关联表简单字段插入------ + Map esFieldData = new LinkedHashMap<>(); + for (FieldItem fieldItem : tableItem.getRelationSelectFieldItems()) { + Object value = esTemplate.getValFromData(config.getEsMapping(), + data, + fieldItem.getFieldName(), + fieldItem.getColumn().getColumnName()); + esFieldData.put(fieldItem.getFieldName(), value); + } + + joinTableSimpleFieldOperation(config, dml, data, tableItem, esFieldData); + } else { + // ------关联子表简单字段插入------ + subTableSimpleFieldOperation(config, dml, data, null, tableItem); + } + } else { + // ------关联子表复杂字段插入 执行全sql更新es------ + jonTableWholeSqlOperation(config, dml, data, null, tableItem); + } + } + } + } + } + + /** + * 更新操作dml + * + * @param config es配置 + * @param dml dml数据 + */ private void update(ESSyncConfig config, Dml dml) { List> dataList = dml.getData(); List> oldList = dml.getOld(); @@ -198,76 +274,6 @@ public class ESSyncService { } } - /** - * 插入dml操作 - * - * @param config es配置 - * @param dml dml数据 - */ - private void insert(ESSyncConfig config, Dml dml) { - List> dataList = dml.getData(); - if (dataList == null || dataList.isEmpty()) { - return; - } - SchemaItem schemaItem = config.getEsMapping().getSchemaItem(); - for (Map data : dataList) { - if (data == null || data.isEmpty()) { - continue; - } - - if (schemaItem.getAliasTableItems().size() == 1 && schemaItem.isAllFieldsSimple()) { - // ------单表 & 所有字段都为简单字段------ - singleTableSimpleFiledInsert(config, dml, data); - } else { - // ------是主表 查询sql来插入------ - if (schemaItem.getMainTable().getTableName().equalsIgnoreCase(dml.getTable())) { - mainTableInsert(config, dml, data); - } - - // 从表的操作 - for (TableItem tableItem : schemaItem.getAliasTableItems().values()) { - if (tableItem.isMain()) { - continue; - } - if (!tableItem.getTableName().equals(dml.getTable())) { - continue; - } - // 关联条件出现在主表查询条件是否为简单字段 - boolean allFieldsSimple = true; - for (FieldItem fieldItem : tableItem.getRelationSelectFieldItems()) { - if (fieldItem.isMethod() || fieldItem.isBinaryOp()) { - allFieldsSimple = false; - break; - } - } - // 所有查询字段均为简单字段 - if (allFieldsSimple) { - // 不是子查询 - if (!tableItem.isSubQuery()) { - // ------关联表简单字段插入------ - Map esFieldData = new LinkedHashMap<>(); - for (FieldItem fieldItem : tableItem.getRelationSelectFieldItems()) { - Object value = esTemplate.getValFromData(config.getEsMapping(), - data, - fieldItem.getFieldName(), - fieldItem.getColumn().getColumnName()); - esFieldData.put(fieldItem.getFieldName(), value); - } - - joinTableSimpleFieldOperation(config, dml, data, tableItem, esFieldData); - } else { - // ------关联子表简单字段插入------ - subTableSimpleFieldOperation(config, dml, data, null, tableItem); - } - } else { - // ------关联子表复杂字段插入 执行全sql更新es------ - jonTableWholeSqlOperation(config, dml, data, null, tableItem); - } - } - } - } - } - /** * 单表简单字段insert * @@ -412,23 +418,7 @@ public class ESSyncService { for (FieldItem fkFieldItem : tableItem.getRelationTableFields().keySet()) { String columnName = fkFieldItem.getColumn().getColumnName(); Object value = esTemplate.getValFromData(mapping, data, fkFieldItem.getFieldName(), columnName); - if (value instanceof String) { - sql.append(tableItem.getAlias()) - .append(".") - .append(columnName) - .append("='") - .append(value) - .append("' ") - .append(" AND "); - } else { - sql.append(tableItem.getAlias()) - .append(".") - .append(columnName) - .append("=") - .append(value) - .append(" ") - .append(" AND "); - } + ESSyncUtil.appendCondition(sql, value, tableItem.getAlias(), columnName); } int len = sql.length(); sql.delete(len - 5, len); @@ -527,23 +517,7 @@ public class ESSyncService { for (FieldItem fkFieldItem : tableItem.getRelationTableFields().keySet()) { String columnName = fkFieldItem.getColumn().getColumnName(); Object value = esTemplate.getValFromData(mapping, data, fkFieldItem.getFieldName(), columnName); - if (value instanceof String) { - sql.append(tableItem.getAlias()) - .append(".") - .append(columnName) - .append("='") - .append(value) - .append("' ") - .append(" AND "); - } else { - sql.append(tableItem.getAlias()) - .append(".") - .append(columnName) - .append("=") - .append(value) - .append(" ") - .append(" AND "); - } + ESSyncUtil.appendCondition(sql, value, tableItem.getAlias(), columnName); } int len = sql.length(); sql.delete(len - 5, len); @@ -602,16 +576,14 @@ public class ESSyncService { BoolQueryBuilder queryBuilder = QueryBuilders.boolQuery(); for (Map.Entry> entry : tableItem.getRelationTableFields().entrySet()) { for (FieldItem fieldItem : entry.getValue()) { - if (fieldItem.getColumnItems().size() == 1) { - Object value = esTemplate - .getValFromRS(mapping, rs, fieldItem.getFieldName(), fieldItem.getFieldName()); - String fieldName = fieldItem.getFieldName(); - // 判断是否是主键 - if (fieldName.equals(mapping.get_id())) { - fieldName = "_id"; - } - queryBuilder.must(QueryBuilders.termsQuery(fieldName, value)); + Object value = esTemplate + .getValFromRS(mapping, rs, fieldItem.getFieldName(), fieldItem.getFieldName()); + String fieldName = fieldItem.getFieldName(); + // 判断是否是主键 + if (fieldName.equals(mapping.get_id())) { + fieldName = "_id"; } + queryBuilder.must(QueryBuilders.termsQuery(fieldName, value)); } } diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESSyncUtil.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESSyncUtil.java index 88f7b496..71331873 100644 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESSyncUtil.java +++ b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/ESSyncUtil.java @@ -19,7 +19,6 @@ import com.alibaba.fastjson.JSON; import com.alibaba.otter.canal.client.adapter.es.config.ESSyncConfig.ESMapping; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem.ColumnItem; -import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem.FieldItem; import com.alibaba.otter.canal.client.adapter.es.config.SchemaItem.TableItem; public class ESSyncUtil { @@ -161,7 +160,7 @@ public class ESSyncUtil { } if (!((String) val).contains(",")) { - logger.error("es type is geo_point, source value not contains ',' spit"); + logger.error("es type is geo_point, source value not contains ',' separator"); return val; } @@ -225,7 +224,7 @@ public class ESSyncUtil { // 拼接condition StringBuilder condition = new StringBuilder(" "); for (ColumnItem idColumn : idColumns) { - Object idVal = data.get(Util.cleanColumn(idColumn.getColumnName())); + Object idVal = data.get(idColumn.getColumnName()); if (mainTable.getAlias() != null) condition.append(mainTable.getAlias()).append("."); condition.append(idColumn.getColumnName()).append("="); if (idVal instanceof String) { @@ -246,6 +245,14 @@ public class ESSyncUtil { return sql + " WHERE " + condition + " "; } + public static void appendCondition(StringBuilder sql, Object value, String owner, String columnName) { + if (value instanceof String) { + sql.append(owner).append(".").append(columnName).append("='").append(value).append("' AND "); + } else { + sql.append(owner).append(".").append(columnName).append("=").append(value).append(" AND "); + } + } + /** * 执行查询sql */ diff --git a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/Util.java b/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/Util.java deleted file mode 100644 index a801b9e6..00000000 --- a/client-adapter/elasticsearch/src/main/java/com/alibaba/otter/canal/client/adapter/es/support/Util.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.alibaba.otter.canal.client.adapter.es.support; - -public class Util { - public static String cleanColumn(String column) { - if (column == null) { - return null; - } - if (column.contains("`")) { - column = column.replaceAll("`", ""); - } - - if (column.contains("'")) { - column = column.replaceAll("'", ""); - } - - return column; - } -} diff --git a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinSub2Test.java b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinSub2Test.java index f67e0b33..469894d3 100644 --- a/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinSub2Test.java +++ b/client-adapter/elasticsearch/src/test/java/com/alibaba/otter/canal/client/adapter/es/test/sync/RoleSyncJoinSub2Test.java @@ -69,7 +69,7 @@ public class RoleSyncJoinSub2Test { List> oldList = new ArrayList<>(); Map old = new LinkedHashMap<>(); oldList.add(old); - old.put("label", "a"); + old.put("label", "v"); dml.setOld(oldList); esAdapter.getEsSyncService().sync(dml);