This is an automated email from the ASF dual-hosted git repository. arvid pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 94da2b587c133621bbdf39d36d55070a64605f56 Author: Fabian Paul <fabianp...@ververica.com> AuthorDate: Mon Aug 16 09:58:23 2021 +0200 [FLINK-23710][connectors/kafka] Move FLIP-143 KafkaSink to org.apache.kafka.connector.kafka.sink --- .../kafka/sink/DefaultKafkaSinkContext.java | 2 +- .../kafka/sink/FlinkKafkaInternalProducer.java | 2 +- .../connectors => connector}/kafka/sink/KafkaCommittable.java | 2 +- .../kafka/sink/KafkaCommittableSerializer.java | 2 +- .../connectors => connector}/kafka/sink/KafkaCommitter.java | 2 +- .../kafka/sink/KafkaRecordSerializationSchema.java | 2 +- .../{streaming/connectors => connector}/kafka/sink/KafkaSink.java | 2 +- .../connectors => connector}/kafka/sink/KafkaSinkBuilder.java | 2 +- .../connectors => connector}/kafka/sink/KafkaTransactionLog.java | 8 ++++---- .../connectors => connector}/kafka/sink/KafkaWriter.java | 2 +- .../connectors => connector}/kafka/sink/KafkaWriterState.java | 2 +- .../kafka/sink/KafkaWriterStateSerializer.java | 2 +- .../kafka/sink/TransactionalIdFactory.java | 2 +- .../kafka/sink/TransactionsToAbortChecker.java | 2 +- .../kafka/table/DynamicKafkaRecordSerializationSchema.java | 8 +++----- .../flink/streaming/connectors/kafka/table/KafkaDynamicSink.java | 4 ++-- .../kafka/sink/KafkaCommittableSerializerTest.java | 2 +- .../connectors => connector}/kafka/sink/KafkaSinkITCase.java | 2 +- .../kafka/sink/KafkaTransactionLogITCase.java | 2 +- .../connectors => connector}/kafka/sink/KafkaWriterITCase.java | 2 +- .../kafka/sink/KafkaWriterStateSerializerTest.java | 2 +- .../kafka/sink/TransactionIdFactoryTest.java | 2 +- .../kafka/sink/TransactionToAbortCheckerTest.java | 2 +- .../connectors/kafka/table/KafkaDynamicTableFactoryTest.java | 2 +- .../kafka/table/UpsertKafkaDynamicTableFactoryTest.java | 2 +- 25 files changed, 31 insertions(+), 33 deletions(-) diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/DefaultKafkaSinkContext.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/DefaultKafkaSinkContext.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/DefaultKafkaSinkContext.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/DefaultKafkaSinkContext.java index 14fc3f2..be4158d 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/DefaultKafkaSinkContext.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/DefaultKafkaSinkContext.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.Producer; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/FlinkKafkaInternalProducer.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/FlinkKafkaInternalProducer.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/FlinkKafkaInternalProducer.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/FlinkKafkaInternalProducer.java index 7a382c7..6866ceb 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/FlinkKafkaInternalProducer.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/FlinkKafkaInternalProducer.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.util.Preconditions; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittable.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittable.java similarity index 97% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittable.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittable.java index 6f49a43..0d5bdb2 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittable.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittable.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.kafka.clients.producer.ProducerConfig; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializer.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializer.java similarity index 97% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializer.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializer.java index f1a6502..7426991 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializer.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializer.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.core.io.SimpleVersionedSerializer; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommitter.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommitter.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommitter.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommitter.java index b411a45..4263656 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommitter.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaCommitter.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.connector.sink.Committer; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaRecordSerializationSchema.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaRecordSerializationSchema.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaRecordSerializationSchema.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaRecordSerializationSchema.java index 4d3b03b..e3ba2fe 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaRecordSerializationSchema.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaRecordSerializationSchema.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.annotation.Internal; import org.apache.flink.annotation.PublicEvolving; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSink.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSink.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSink.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSink.java index fd3c2d7..7e6b033 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSink.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSink.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.connector.sink.Committer; import org.apache.flink.api.connector.sink.GlobalCommitter; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkBuilder.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSinkBuilder.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkBuilder.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSinkBuilder.java index ab2aca5..4428e89 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkBuilder.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaSinkBuilder.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.common.ExecutionConfig; import org.apache.flink.api.java.ClosureCleaner; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLog.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLog.java similarity index 96% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLog.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLog.java index 9a26359..4e5926e 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLog.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLog.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList; import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableSet; @@ -46,9 +46,9 @@ import java.util.UUID; import java.util.function.Consumer; import java.util.stream.Collectors; -import static org.apache.flink.streaming.connectors.kafka.sink.KafkaTransactionLog.TransactionState.CompleteAbort; -import static org.apache.flink.streaming.connectors.kafka.sink.KafkaTransactionLog.TransactionState.CompleteCommit; -import static org.apache.flink.streaming.connectors.kafka.sink.KafkaTransactionLog.TransactionState.Dead; +import static org.apache.flink.connector.kafka.sink.KafkaTransactionLog.TransactionState.CompleteAbort; +import static org.apache.flink.connector.kafka.sink.KafkaTransactionLog.TransactionState.CompleteCommit; +import static org.apache.flink.connector.kafka.sink.KafkaTransactionLog.TransactionState.Dead; import static org.apache.flink.util.Preconditions.checkNotNull; import static org.apache.kafka.common.internals.Topic.TRANSACTION_STATE_TOPIC_NAME; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriter.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriter.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriter.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriter.java index c05dc6b..06c32e6 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriter.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriter.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.common.serialization.SerializationSchema; import org.apache.flink.api.connector.sink.Sink; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterState.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterState.java similarity index 97% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterState.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterState.java index 5868c1a..b5187b6 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterState.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterState.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import java.util.Objects; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializer.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializer.java similarity index 97% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializer.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializer.java index f2487fe..609b899 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializer.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializer.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.core.io.SimpleVersionedSerializer; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionalIdFactory.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionalIdFactory.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionalIdFactory.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionalIdFactory.java index b9c5c5e..96744b6 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionalIdFactory.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionalIdFactory.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.shaded.guava30.com.google.common.base.Splitter; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionsToAbortChecker.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionsToAbortChecker.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionsToAbortChecker.java rename to flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionsToAbortChecker.java index 1e21002..2d79b35 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionsToAbortChecker.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/sink/TransactionsToAbortChecker.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import java.util.ArrayList; import java.util.Collections; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaRecordSerializationSchema.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaRecordSerializationSchema.java index 45ae5a8..7908ade 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaRecordSerializationSchema.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/DynamicKafkaRecordSerializationSchema.java @@ -18,8 +18,9 @@ package org.apache.flink.streaming.connectors.kafka.table; import org.apache.flink.api.common.serialization.SerializationSchema; +import org.apache.flink.connector.kafka.sink.KafkaRecordSerializationSchema; +import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkKafkaPartitioner; -import org.apache.flink.streaming.connectors.kafka.sink.KafkaRecordSerializationSchema; import org.apache.flink.table.data.GenericRowData; import org.apache.flink.table.data.RowData; import org.apache.flink.types.RowKind; @@ -31,10 +32,7 @@ import javax.annotation.Nullable; import static org.apache.flink.util.Preconditions.checkNotNull; -/** - * SerializationSchema used by {@link KafkaDynamicSink} to configure a {@link - * org.apache.flink.streaming.connectors.kafka.sink.KafkaSink}. - */ +/** SerializationSchema used by {@link KafkaDynamicSink} to configure a {@link KafkaSink}. */ class DynamicKafkaRecordSerializationSchema implements KafkaRecordSerializationSchema<RowData> { private final String topic; diff --git a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSink.java b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSink.java index 9fd5af8..1fd9caf 100644 --- a/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSink.java +++ b/flink-connectors/flink-connector-kafka/src/main/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicSink.java @@ -24,11 +24,11 @@ import org.apache.flink.api.common.serialization.SerializationSchema; import org.apache.flink.api.common.typeinfo.TypeInformation; import org.apache.flink.api.common.typeutils.TypeSerializer; import org.apache.flink.connector.base.DeliveryGuarantee; +import org.apache.flink.connector.kafka.sink.KafkaSink; +import org.apache.flink.connector.kafka.sink.KafkaSinkBuilder; import org.apache.flink.streaming.api.datastream.DataStreamSink; import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer; import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkKafkaPartitioner; -import org.apache.flink.streaming.connectors.kafka.sink.KafkaSink; -import org.apache.flink.streaming.connectors.kafka.sink.KafkaSinkBuilder; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.connector.ChangelogMode; import org.apache.flink.table.connector.format.EncodingFormat; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializerTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializerTest.java similarity index 96% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializerTest.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializerTest.java index bc51ad6..a55bb4f 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaCommittableSerializerTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaCommittableSerializerTest.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.junit.Test; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkITCase.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkITCase.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java index cf7cf4c..3eb2c18 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaSinkITCase.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.common.functions.MapFunction; import org.apache.flink.api.common.restartstrategy.RestartStrategies; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLogITCase.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLogITCase.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLogITCase.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLogITCase.java index eed877b..13820fb 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaTransactionLogITCase.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaTransactionLogITCase.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableList; import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableMap; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterITCase.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterITCase.java similarity index 99% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterITCase.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterITCase.java index f9c338e..4cb2ae1 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterITCase.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterITCase.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.api.common.operators.MailboxExecutor; import org.apache.flink.api.common.serialization.SerializationSchema; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializerTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializerTest.java similarity index 96% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializerTest.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializerTest.java index 3c52895..dea374e 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/KafkaWriterStateSerializerTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaWriterStateSerializerTest.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.junit.Test; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionIdFactoryTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionIdFactoryTest.java similarity index 97% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionIdFactoryTest.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionIdFactoryTest.java index fc36baa..fbe9b17 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionIdFactoryTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionIdFactoryTest.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.junit.Test; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionToAbortCheckerTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionToAbortCheckerTest.java similarity index 98% rename from flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionToAbortCheckerTest.java rename to flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionToAbortCheckerTest.java index 4891147..d622661 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/sink/TransactionToAbortCheckerTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/TransactionToAbortCheckerTest.java @@ -15,7 +15,7 @@ * limitations under the License. */ -package org.apache.flink.streaming.connectors.kafka.sink; +package org.apache.flink.connector.kafka.sink; import org.apache.flink.shaded.guava30.com.google.common.collect.ImmutableMap; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java index 7508749..01af4b0 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/KafkaDynamicTableFactoryTest.java @@ -23,6 +23,7 @@ import org.apache.flink.api.common.serialization.SerializationSchema; import org.apache.flink.api.connector.sink.Sink; import org.apache.flink.api.dag.Transformation; import org.apache.flink.connector.base.DeliveryGuarantee; +import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumState; import org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit; @@ -37,7 +38,6 @@ import org.apache.flink.streaming.connectors.kafka.config.StartupMode; import org.apache.flink.streaming.connectors.kafka.internals.KafkaTopicPartition; import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkFixedPartitioner; import org.apache.flink.streaming.connectors.kafka.partitioner.FlinkKafkaPartitioner; -import org.apache.flink.streaming.connectors.kafka.sink.KafkaSink; import org.apache.flink.streaming.connectors.kafka.table.KafkaConnectorOptions.ScanStartupMode; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.api.ValidationException; diff --git a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/UpsertKafkaDynamicTableFactoryTest.java b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/UpsertKafkaDynamicTableFactoryTest.java index 5138ebc..16cd90f 100644 --- a/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/UpsertKafkaDynamicTableFactoryTest.java +++ b/flink-connectors/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/table/UpsertKafkaDynamicTableFactoryTest.java @@ -23,6 +23,7 @@ import org.apache.flink.api.common.serialization.SerializationSchema; import org.apache.flink.api.connector.sink.Sink; import org.apache.flink.api.dag.Transformation; import org.apache.flink.connector.base.DeliveryGuarantee; +import org.apache.flink.connector.kafka.sink.KafkaSink; import org.apache.flink.connector.kafka.source.KafkaSource; import org.apache.flink.connector.kafka.source.enumerator.KafkaSourceEnumState; import org.apache.flink.connector.kafka.source.split.KafkaPartitionSplit; @@ -34,7 +35,6 @@ import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.operators.StreamOperatorFactory; import org.apache.flink.streaming.api.transformations.SourceTransformation; import org.apache.flink.streaming.connectors.kafka.config.StartupMode; -import org.apache.flink.streaming.connectors.kafka.sink.KafkaSink; import org.apache.flink.streaming.runtime.operators.sink.SinkOperatorFactory; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.api.ValidationException;