Skip to content

Commit 214562d

Browse files
author
dapeng
committed
测试日志
1 parent 174f70f commit 214562d

2 files changed

Lines changed: 2 additions & 0 deletions

File tree

kafka-base/kafka-base-sink/src/main/java/com/dtstack/flink/sql/sink/kafka/CustomerFlinkPartition.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
public class CustomerFlinkPartition<T> extends FlinkFixedPartitioner<T> {
1010
@Override
1111
public int partition(T record, byte[] key, byte[] value, String targetTopic, int[] partitions) {
12+
System.out.println("key="+new String(key));
1213
Preconditions.checkArgument(partitions != null && partitions.length > 0, "Partitions of the target topic is empty.");
1314
if(key == null){
1415
Random random = new Random();

kafka-base/kafka-base-sink/src/main/java/com/dtstack/flink/sql/sink/kafka/CustomerKeyedSerializationSchema.java

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ public CustomerKeyedSerializationSchema(SerializationMetricWrapper serialization
2323
}
2424

2525
public byte[] serializeKey(Row element) {
26+
System.out.println("element = " + element+"|partitionKeys=" + partitionKeys);
2627
if(partitionKeys == null || partitionKeys.length <=0){
2728
return null;
2829
}

0 commit comments

Comments
 (0)