Skip to content

Commit c0c7a23

Browse files
author
dapeng
committed
remove logger
1 parent 286a796 commit c0c7a23

1 file changed

Lines changed: 2 additions & 10 deletions

File tree

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

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -2,23 +2,15 @@
22

33

44
import com.dtstack.flink.sql.format.SerializationMetricWrapper;
5-
import org.apache.flink.api.common.functions.RuntimeContext;
65
import org.apache.flink.api.common.serialization.SerializationSchema;
7-
import org.apache.flink.formats.json.JsonRowSchemaConverter;
86
import org.apache.flink.formats.json.JsonRowSerializationSchema;
9-
import org.apache.flink.metrics.Counter;
10-
import org.apache.flink.metrics.Meter;
117
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper;
8+
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ObjectNode;
129
import org.apache.flink.streaming.util.serialization.KeyedSerializationSchema;
1310
import org.apache.flink.types.Row;
14-
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.node.ObjectNode;
15-
import org.slf4j.Logger;
16-
import org.slf4j.LoggerFactory;
1711

1812
public class CustomerKeyedSerializationSchema implements KeyedSerializationSchema<Row> {
1913

20-
private Logger log = LoggerFactory.getLogger(getClass());
21-
2214
private static final long serialVersionUID = 1L;
2315
private final SerializationMetricWrapper serializationMetricWrapper;
2416
private String[] partitionKeys;
@@ -61,7 +53,7 @@ private byte[] serializeJsonKey(JsonRowSerializationSchema jsonRowSerializationS
6153
}
6254
return sb.toString().getBytes();
6355
}catch (Exception e){
64-
log.error("serializeJsonKey error", e);
56+
6557
}
6658
return null;
6759

0 commit comments

Comments
 (0)