Commit 20038450 authored by 375138141's avatar 375138141

kingbase

parent ffabf113
...@@ -81,8 +81,8 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple ...@@ -81,8 +81,8 @@ public class MysqlDataTransferFunction extends ProcessFunction<JSONObject, Tuple
} }
String groupKey = groupKeyBuilder.toString(); String groupKey = groupKeyBuilder.toString();
//添加分表参数 //添加分表参数(shardingRule 位于消息顶层,由 dsk-flink-upgrade 注入)
String shardingRule = dataObj.getString("shardingRule"); String shardingRule = value.getString("shardingRule");
if (StrUtil.isNotBlank(shardingRule)) { if (StrUtil.isNotBlank(shardingRule)) {
Map<String,Object> map = JSON.parseObject(shardingRule, Map.class); Map<String,Object> map = JSON.parseObject(shardingRule, Map.class);
String strategy = MapUtil.getStr(map, "strategy"); String strategy = MapUtil.getStr(map, "strategy");
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment