This is an automated email from the ASF dual-hosted git repository. dwysakowicz pushed a change to branch master in repository https://gitbox.apache.org/repos/asf/flink.git.
from ad46ca3 [FLINK-17416][e2e][k8s][hotfix] Disable failing test add c20374b [FLINK-17306] Added open to DeserializationSchema add d97d75a [FLINK-17306] Add open to KafkaDeserializationSchema add 8b128c9 [FLINK-17306] Add open to KinesisDeserializationSchema add ff07c3a [FLINK-17306] Add open to PubSubDeserializationSchema add e2973f1 [FLINK-17306] Call open of DeserializationSchema in RMQ add cf68178 [FLINK-17306] Add open to SerializationSchema add 1849f3a [FLINK-17306] Add open to KafkaSerializationSchema add 6aa50ed [FLINK-17306] Add open to KinesisSerializationSchema add 64f4a43 [FLINK-17306] Call open of SerializationSchema in PubSub sink add 22334ff [FLINK-17306] Call open of SerializationSchema in RMQ sink No new revisions were added by this update. Summary of changes: .../gcp/pubsub/DeserializationSchemaWrapper.java | 5 ++ .../connectors/gcp/pubsub/PubSubSink.java | 2 + .../connectors/gcp/pubsub/PubSubSource.java | 1 + .../pubsub/common/PubSubDeserializationSchema.java | 13 +++++ .../pubsub/DeserializationSchemaWrapperTest.java | 12 ++++ .../connectors/gcp/pubsub/PubSubSourceTest.java | 17 ++++++ .../connectors/kafka/FlinkKafkaProducer011.java | 5 ++ .../connectors/kafka/FlinkKafkaConsumerBase.java | 2 + .../connectors/kafka/FlinkKafkaProducerBase.java | 7 ++- .../kafka/KafkaDeserializationSchema.java | 13 +++++ .../connectors/kafka/KafkaSerializationSchema.java | 13 +++++ .../KafkaDeserializationSchemaWrapper.java | 5 ++ .../internals/KeyedSerializationSchemaWrapper.java | 4 ++ .../kafka/FlinkKafkaConsumerBaseTest.java | 40 +++++++++++-- .../connectors/kafka/FlinkKafkaProducer.java | 4 ++ .../connectors/kinesis/FlinkKinesisProducer.java | 8 +++ .../internals/DynamoDBStreamsDataFetcher.java | 8 ++- .../kinesis/internals/KinesisDataFetcher.java | 31 +++++++--- .../kinesis/internals/ShardConsumer.java | 5 +- .../KinesisDeserializationSchema.java | 12 ++++ .../KinesisDeserializationSchemaWrapper.java | 5 ++ .../serialization/KinesisSerializationSchema.java | 13 +++++ .../kinesis/FlinkKinesisConsumerTest.java | 23 +++++++- .../kinesis/FlinkKinesisProducerTest.java | 27 ++++++++- .../kinesis/internals/ShardConsumerTest.java | 32 ++++++++--- ...inesisDataFetcherForShardConsumerException.java | 5 -- .../streaming/connectors/rabbitmq/RMQSink.java | 2 + .../streaming/connectors/rabbitmq/RMQSource.java | 1 + .../streaming/connectors/rabbitmq/RMQSinkTest.java | 18 ++++++ .../connectors/rabbitmq/RMQSourceTest.java | 30 +++++++++- .../serialization/DeserializationSchema.java | 35 +++++++++++ .../common/serialization/SerializationSchema.java | 35 +++++++++++ .../streaming/util/MockDeserializationSchema.java | 67 ++++++++++++++++++++++ .../streaming/util/MockSerializationSchema.java | 54 +++++++++++++++++ 34 files changed, 518 insertions(+), 36 deletions(-) create mode 100644 flink-streaming-java/src/test/java/org/apache/flink/streaming/util/MockDeserializationSchema.java create mode 100644 flink-streaming-java/src/test/java/org/apache/flink/streaming/util/MockSerializationSchema.java