From 6de9f7efa730afd9cf0f5607376773b9e453318c Mon Sep 17 00:00:00 2001 From: machey Date: Tue, 6 Nov 2018 22:49:34 +0800 Subject: [PATCH 1/6] =?UTF-8?q?=E4=BF=AE=E6=94=B9=20adapter=20release=20?= =?UTF-8?q?=E6=89=93=E5=8C=85=E5=90=8D=E7=A7=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client-adapter/launcher/pom.xml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/client-adapter/launcher/pom.xml b/client-adapter/launcher/pom.xml index 6b5f71dd..6c39f828 100644 --- a/client-adapter/launcher/pom.xml +++ b/client-adapter/launcher/pom.xml @@ -216,7 +216,7 @@ ${basedir}/src/main/assembly/release.xml - ${project.artifactId}-${project.version} + canal.adapter-${project.version} ${project.basedir}/../../target From 3090557eeb8dbc9712bb8f5aa237ec249e957cc3 Mon Sep 17 00:00:00 2001 From: machey Date: Tue, 6 Nov 2018 23:16:11 +0800 Subject: [PATCH 2/6] =?UTF-8?q?client=20=E6=8E=92=E9=99=A4jsr305=E4=B8=8D?= =?UTF-8?q?=E6=89=93=E8=BF=9Bfat=20jar?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- client/pom.xml | 1 + 1 file changed, 1 insertion(+) diff --git a/client/pom.xml b/client/pom.xml index c12447c9..aa31390a 100644 --- a/client/pom.xml +++ b/client/pom.xml @@ -149,6 +149,7 @@ commons-logging:commons-logging aopalliance:aopalliance com.google.protobuf:protobuf-java + com.google.code.findbugs:jsr305 io.netty:* junit:junit From c27d5fa16274207176f9a638ea92e3afca118189 Mon Sep 17 00:00:00 2001 From: machey Date: Wed, 7 Nov 2018 02:04:08 +0800 Subject: [PATCH 3/6] =?UTF-8?q?ETL=E6=95=B4=E4=B8=AAdestination=EF=BC=8C?= =?UTF-8?q?=E5=90=8C=E6=AD=A5=E5=BC=80=E5=85=B3=E6=97=A0=E6=95=88=E7=9A=84?= =?UTF-8?q?=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/launcher/rest/CommonRest.java | 18 ++++++++++++++---- 1 file changed, 14 insertions(+), 4 deletions(-) diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java index d2c163fb..6f09c33b 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java @@ -67,10 +67,18 @@ public class CommonRest { try { OuterAdapter adapter = loader.getExtension(type); String destination = adapter.getDestination(task); - Boolean oriSwithcStatus = null; + Boolean oriSwitchStatus; if (destination != null) { - oriSwithcStatus = syncSwitch.status(destination); - syncSwitch.off(destination); + oriSwitchStatus = syncSwitch.status(destination); + if (oriSwitchStatus != null && oriSwitchStatus) { + syncSwitch.off(destination); + } + } else { + // task可能为destination,直接锁task + oriSwitchStatus = syncSwitch.status(task); + if (oriSwitchStatus != null && oriSwitchStatus) { + syncSwitch.off(task); + } } try { List paramArr = null; @@ -80,8 +88,10 @@ public class CommonRest { } return adapter.etl(task, paramArr); } finally { - if (destination != null && oriSwithcStatus != null && oriSwithcStatus) { + if (destination != null && oriSwitchStatus != null && oriSwitchStatus) { syncSwitch.on(destination); + } else if (destination == null && oriSwitchStatus != null && oriSwitchStatus) { + syncSwitch.on(task); } } } finally { From 202557630c027ffb23bac383bd87d409b5e751b4 Mon Sep 17 00:00:00 2001 From: mcy Date: Wed, 7 Nov 2018 09:05:47 +0800 Subject: [PATCH 4/6] =?UTF-8?q?=E8=B0=83=E6=95=B4etl=E9=94=81=E7=9A=84?= =?UTF-8?q?=E7=B2=92=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../otter/canal/adapter/launcher/rest/CommonRest.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java index 6f09c33b..0a6d00f4 100644 --- a/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java +++ b/client-adapter/launcher/src/main/java/com/alibaba/otter/canal/adapter/launcher/rest/CommonRest.java @@ -56,8 +56,11 @@ public class CommonRest { @PostMapping("/etl/{type}/{task}") public EtlResult etl(@PathVariable String type, @PathVariable String task, @RequestParam(name = "params", required = false) String params) { + OuterAdapter adapter = loader.getExtension(type); + String destination = adapter.getDestination(task); + String lockKey = destination == null ? task : destination; - boolean locked = etlLock.tryLock(ETL_LOCK_ZK_NODE + type + "-" + task); + boolean locked = etlLock.tryLock(ETL_LOCK_ZK_NODE + type + "-" + lockKey); if (!locked) { EtlResult result = new EtlResult(); result.setSucceeded(false); @@ -65,8 +68,7 @@ public class CommonRest { return result; } try { - OuterAdapter adapter = loader.getExtension(type); - String destination = adapter.getDestination(task); + Boolean oriSwitchStatus; if (destination != null) { oriSwitchStatus = syncSwitch.status(destination); @@ -95,7 +97,7 @@ public class CommonRest { } } } finally { - etlLock.unlock(ETL_LOCK_ZK_NODE + type + "-" + task); + etlLock.unlock(ETL_LOCK_ZK_NODE + type + "-" + lockKey); } } From bc298bedd533bde6247cba6b59207d70f266a4b5 Mon Sep 17 00:00:00 2001 From: mcy Date: Wed, 7 Nov 2018 09:32:52 +0800 Subject: [PATCH 5/6] =?UTF-8?q?=E6=94=AF=E6=8C=81hbase=20rowkey=E6=8C=87?= =?UTF-8?q?=E5=AE=9A=E6=95=B0=E5=AD=97=E9=95=BF=E5=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../adapter/hbase/config/MappingConfig.java | 20 +++++- .../hbase/service/HbaseEtlService.java | 67 +++++++++++++++---- .../hbase/service/HbaseSyncService.java | 43 ++++++++++-- .../main/resources/hbase/mytest_person2.yml | 14 ++-- 4 files changed, 118 insertions(+), 26 deletions(-) diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java index c380b65c..326971ac 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/config/MappingConfig.java @@ -76,6 +76,7 @@ public class MappingConfig { public static class ColumnItem { private boolean isRowKey = false; + private Integer rowKeyLen; private String column; private String family; private String qualifier; @@ -89,6 +90,14 @@ public class MappingConfig { isRowKey = rowKey; } + public Integer getRowKeyLen() { + return rowKeyLen; + } + + public void setRowKeyLen(Integer rowKeyLen) { + this.rowKeyLen = rowKeyLen; + } + public String getColumn() { return column; } @@ -264,7 +273,16 @@ public class MappingConfig { ColumnItem columnItem = new ColumnItem(); columnItem.setColumn(columnField.getKey()); columnItem.setType(type); - if ("rowKey".equalsIgnoreCase(field)) { + if (field != null && field.toUpperCase().startsWith("ROWKEY")) { + int idx = field.toUpperCase().indexOf("LEN:"); + if (idx > -1) { + String len = field.substring(idx + 4); + try { + columnItem.setRowKeyLen(Integer.parseInt(len)); + } catch (Exception e) { + // ignore + } + } columnItem.setRowKey(true); rowKeyColumn = columnItem; } else { diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java index 7c57a9ab..912baf08 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseEtlService.java @@ -199,8 +199,8 @@ public class HbaseEtlService { } else { sqlFinal = sql + " LIMIT " + offset + "," + cnt; } - Future future = executor - .submit(() -> executeSqlImport(ds, sqlFinal, hbaseMapping, hbaseTemplate, successCount, errMsg)); + Future future = executor.submit( + () -> executeSqlImport(ds, sqlFinal, hbaseMapping, hbaseTemplate, successCount, errMsg)); futures.add(future); } @@ -213,8 +213,8 @@ public class HbaseEtlService { executeSqlImport(ds, sql, hbaseMapping, hbaseTemplate, successCount, errMsg); } - logger.info( - hbaseMapping.getHbaseTable() + " etl completed in: " + (System.currentTimeMillis() - start) / 1000 + "s!"); + logger.info(hbaseMapping.getHbaseTable() + " etl completed in: " + + (System.currentTimeMillis() - start) / 1000 + "s!"); etlResult.setResultMessage("导入HBase表 " + hbaseMapping.getHbaseTable() + " 数据:" + successCount.get() + " 条"); } catch (Exception e) { @@ -332,18 +332,49 @@ public class HbaseEtlService { byte[] valBytes = Bytes.toBytes(val.toString()); if (columnItem.isRowKey()) { - row.setRowKey(valBytes); + if (columnItem.getRowKeyLen() != null) { + valBytes = Bytes.toBytes(limitLenNum(columnItem.getRowKeyLen(), val)); + row.setRowKey(valBytes); + } else { + row.setRowKey(valBytes); + } } else { row.addCell(columnItem.getFamily(), columnItem.getQualifier(), valBytes); } } else { - PhType phType = PhType.getType(columnItem.getType()); - if (columnItem.isRowKey()) { - row.setRowKey(PhTypeUtil.toBytes(val, phType)); - } else { - row.addCell(columnItem.getFamily(), - columnItem.getQualifier(), - PhTypeUtil.toBytes(val, phType)); + if (MappingConfig.Mode.STRING == hbaseMapping.getMode()) { + byte[] valBytes = Bytes.toBytes(val.toString()); + if (columnItem.isRowKey()) { + if (columnItem.getRowKeyLen() != null) { + valBytes = Bytes.toBytes(limitLenNum(columnItem.getRowKeyLen(), val)); + } + row.setRowKey(valBytes); + } else { + row.addCell(columnItem.getFamily(), columnItem.getQualifier(), valBytes); + } + } else if (MappingConfig.Mode.NATIVE == hbaseMapping.getMode()) { + Type type = Type.getType(columnItem.getType()); + if (columnItem.isRowKey()) { + if (columnItem.getRowKeyLen() != null) { + String v = limitLenNum(columnItem.getRowKeyLen(), val); + row.setRowKey(Bytes.toBytes(v)); + } else { + row.setRowKey(TypeUtil.toBytes(val, type)); + } + } else { + row.addCell(columnItem.getFamily(), + columnItem.getQualifier(), + TypeUtil.toBytes(val, type)); + } + } else if (MappingConfig.Mode.PHOENIX == hbaseMapping.getMode()) { + PhType phType = PhType.getType(columnItem.getType()); + if (columnItem.isRowKey()) { + row.setRowKey(PhTypeUtil.toBytes(val, phType)); + } else { + row.addCell(columnItem.getFamily(), + columnItem.getQualifier(), + PhTypeUtil.toBytes(val, phType)); + } } } } @@ -382,4 +413,16 @@ public class HbaseEtlService { return false; } } + + private static String limitLenNum(int len, Object val) { + if (val == null) { + return null; + } + if (val instanceof Number) { + return String.format("%0" + len + "d", (Number) ((Number) val).longValue()); + } else if (val instanceof String) { + return String.format("%0" + len + "d", Long.parseLong((String) val)); + } + return null; + } } diff --git a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java index 420640af..2a242093 100644 --- a/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java +++ b/client-adapter/hbase/src/main/java/com/alibaba/otter/canal/client/adapter/hbase/service/HbaseSyncService.java @@ -132,7 +132,21 @@ public class HbaseSyncService { } } else { if (columnItem.isRowKey()) { - // row.put("rowKey", bytes); + if (columnItem.getRowKeyLen() != null && entry.getValue() != null) { + if (entry.getValue() instanceof Number) { + String v = String.format("%0" + columnItem.getRowKeyLen() + "d", + ((Number) entry.getValue()).longValue()); + bytes = Bytes.toBytes(v); + } else { + try { + String v = String.format("%0" + columnItem.getRowKeyLen() + "d", + Integer.parseInt((String) entry.getValue())); + bytes = Bytes.toBytes(v); + } catch (Exception e) { + // ignore + } + } + } hRow.setRowKey(bytes); } else { hRow.addCell(columnItem.getFamily(), columnItem.getQualifier(), bytes); @@ -146,8 +160,8 @@ public class HbaseSyncService { /** * 更新操作 * - * @param config - * @param dml + * @param config 配置对象 + * @param dml dml对象 */ private void update(MappingConfig config, Dml dml) { List> data = dml.getData(); @@ -192,7 +206,7 @@ public class HbaseSyncService { Map rowKey = data.get(0); rowKeyBytes = typeConvert(null, hbaseMapping, rowKey.values().iterator().next()); } else { - rowKeyBytes = typeConvert(rowKeyColumn, hbaseMapping, r.get(rowKeyColumn.getColumn())); + rowKeyBytes = getRowKeyBytes(hbaseMapping, rowKeyColumn, r); } if (rowKeyBytes == null) throw new RuntimeException("rowKey值为空"); @@ -277,8 +291,7 @@ public class HbaseSyncService { Map rowKey = data.get(0); rowKeyBytes = typeConvert(null, hbaseMapping, rowKey.values().iterator().next()); } else { - Object val = r.get(rowKeyColumn.getColumn()); - rowKeyBytes = typeConvert(rowKeyColumn, hbaseMapping, val); + rowKeyBytes = getRowKeyBytes(hbaseMapping, rowKeyColumn, r); } if (rowKeyBytes == null) throw new RuntimeException("rowKey值为空"); rowKeys.add(rowKeyBytes); @@ -424,4 +437,22 @@ public class HbaseSyncService { return rowKeyValue.toString(); } + private static byte[] getRowKeyBytes(MappingConfig.HbaseMapping hbaseMapping, MappingConfig.ColumnItem rowKeyColumn, + Map rowData) { + Object val = rowData.get(rowKeyColumn.getColumn()); + String v = null; + if (rowKeyColumn.getRowKeyLen() != null) { + if (val instanceof Number) { + v = String.format("%0" + rowKeyColumn.getRowKeyLen() + "d", (Number) ((Number) val).longValue()); + } else if (val instanceof String) { + v = String.format("%0" + rowKeyColumn.getRowKeyLen() + "d", Long.parseLong((String) val)); + } + } + if (v != null) { + return Bytes.toBytes(v); + } else { + return typeConvert(rowKeyColumn, hbaseMapping, val); + } + } + } diff --git a/client-adapter/hbase/src/main/resources/hbase/mytest_person2.yml b/client-adapter/hbase/src/main/resources/hbase/mytest_person2.yml index cbfe6bb1..0eebad47 100644 --- a/client-adapter/hbase/src/main/resources/hbase/mytest_person2.yml +++ b/client-adapter/hbase/src/main/resources/hbase/mytest_person2.yml @@ -1,7 +1,7 @@ dataSourceKey: defaultDS destination: example hbaseMapping: - mode: PHOENIX #NATIVE #STRING + mode: STRING #NATIVE #PHOENIX database: mytest # 数据库名 table: person2 # 数据库表名 hbaseTable: MYTEST.PERSON2 # HBase表名 @@ -11,14 +11,14 @@ hbaseMapping: #rowKey: id,type # 复合字段rowKey不能和columns中的rowKey重复 columns: # 数据库字段:HBase对应字段 - id: ROWKEY$UNSIGNED_LONG + id: ROWKEY LEN:15 name: NAME email: EMAIL - type: $DECIMAL - c_time: C_TIME$UNSIGNED_TIMESTAMP - birthday: BIRTHDAY$DATE - excludeColumns: - - lat # 忽略字段 + type: + c_time: C_TIME + birthday: BIRTHDAY +# excludeColumns: +# - lat # 忽略字段 # -- NATIVE类型 # $DEFAULT From b71e2991610febd5d19d9e3a303d3c9fe667ce32 Mon Sep 17 00:00:00 2001 From: mcy Date: Wed, 7 Nov 2018 11:37:16 +0800 Subject: [PATCH 6/6] =?UTF-8?q?kafka=E5=8D=95=E4=B8=AApartition=E7=9A=84?= =?UTF-8?q?=E5=8F=91=E9=80=81=E6=9C=AA=E6=94=B9=E4=B8=BA=E5=90=8C=E6=AD=A5?= =?UTF-8?q?=E5=8F=91=E9=80=81=E7=9A=84bug=20fix?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java index 8576835f..a4c9eaa8 100644 --- a/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java +++ b/server/src/main/java/com/alibaba/otter/canal/kafka/CanalKafkaProducer.java @@ -105,7 +105,7 @@ public class CanalKafkaProducer implements CanalMQProducer { canalDestination.getPartition(), null, JSON.toJSONString(flatMessage)); - producer2.send(record); + producer2.send(record).get(); } catch (Exception e) { logger.error(e.getMessage(), e); // producer.abortTransaction();