diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java index 77e22415..ebcb4ec0 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/JdbcTypeUtil.java @@ -76,7 +76,7 @@ public class JdbcTypeUtil { || "TEXT".equalsIgnoreCase(columnType) || "TINYTEXT".equalsIgnoreCase(columnType); } - public static Object typeConvert(String columnName, String value, int sqlType, String mysqlType) { + public static Object typeConvert(String tableName ,String columnName, String value, int sqlType, String mysqlType) { if (value == null || (value.equals("") && !(isText(mysqlType) || sqlType == Types.CHAR || sqlType == Types.VARCHAR || sqlType == Types.LONGVARCHAR))) { return null; @@ -161,7 +161,7 @@ public class JdbcTypeUtil { } return res; } catch (Exception e) { - logger.error("table: {} column: {}, failed convert type {} to {}", columnName, value, sqlType); + logger.error("table: {} column: {}, failed convert type {} to {}", tableName, columnName, value, sqlType); return value; } } diff --git a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java index 51e6a1b4..1e24d9ad 100644 --- a/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java +++ b/client-adapter/common/src/main/java/com/alibaba/otter/canal/client/adapter/support/MessageUtil.java @@ -77,7 +77,7 @@ public class MessageUtil { } } row.put(column.getName(), - JdbcTypeUtil.typeConvert(column.getName(), + JdbcTypeUtil.typeConvert(dml.getTable(),column.getName(), column.getValue(), column.getSqlType(), column.getMysqlType())); @@ -95,7 +95,7 @@ public class MessageUtil { for (CanalEntry.Column column : rowData.getBeforeColumnsList()) { if (updateSet.contains(column.getName())) { rowOld.put(column.getName(), - JdbcTypeUtil.typeConvert(column.getName(), + JdbcTypeUtil.typeConvert(dml.getTable(),column.getName(), column.getValue(), column.getSqlType(), column.getMysqlType())); @@ -153,16 +153,16 @@ public class MessageUtil { // } List> data = flatMessage.getData(); if (data != null) { - dml.setData(changeRows(data, flatMessage.getSqlType(), flatMessage.getMysqlType())); + dml.setData(changeRows(dml.getTable(), data, flatMessage.getSqlType(), flatMessage.getMysqlType())); } List> old = flatMessage.getOld(); if (old != null) { - dml.setOld(changeRows(old, flatMessage.getSqlType(), flatMessage.getMysqlType())); + dml.setOld(changeRows(dml.getTable(), old, flatMessage.getSqlType(), flatMessage.getMysqlType())); } return dml; } - private static List> changeRows(List> rows, Map sqlTypes, + private static List> changeRows(String table, List> rows, Map sqlTypes, Map mysqlTypes) { List> result = new ArrayList<>(); for (Map row : rows) { @@ -181,7 +181,7 @@ public class MessageUtil { continue; } - Object finalValue = JdbcTypeUtil.typeConvert(columnName, columnValue, sqlType, mysqlType); + Object finalValue = JdbcTypeUtil.typeConvert(table, columnName, columnValue, sqlType, mysqlType); resultRow.put(columnName, finalValue); } result.add(resultRow);