Skip to content

Commit b1b7af8

Browse files
committed
* Splits Elastic module into elastic-common and elastic8
* Includes refactoring of all connectors to ensure that Maps are passed around and util.Maps are only used on entry and exit from scala code and converted soon after * Adding Opensearch module * Adding Opensearch unit tests * Opensearch SSL Test - needs completing * Replace Elastic6+7 with Elastic8 Removing deleted modules Source package renamer Package rename Dir rename
1 parent 7d8f2dc commit b1b7af8

133 files changed

Lines changed: 3034 additions & 4006 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

build.sbt

Lines changed: 39 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
import Dependencies.globalExcludeDeps
22
import Dependencies.gson
3+
import Dependencies.bouncyCastle
4+
35
import Settings.*
46
import sbt.Keys.libraryDependencies
57
import sbt.*
@@ -17,8 +19,9 @@ lazy val subProjects: Seq[Project] = Seq(
1719
`azure-documentdb`,
1820
`azure-datalake`,
1921
cassandra,
20-
elastic6,
21-
elastic7,
22+
`elastic-common`,
23+
opensearch,
24+
elastic8,
2225
ftp,
2326
`gcp-storage`,
2427
http,
@@ -201,17 +204,16 @@ lazy val cassandra = (project in file("kafka-connect-cassandra"))
201204
.configureFunctionalTests()
202205
.enablePlugins(PackPlugin)
203206

204-
lazy val elastic6 = (project in file("kafka-connect-elastic6"))
207+
lazy val `elastic-common` = (project in file("kafka-connect-elastic-common"))
205208
.dependsOn(common)
206209
.dependsOn(`test-common` % "fun->compile")
207210
.settings(
208211
settings ++
209212
Seq(
210-
name := "kafka-connect-elastic6",
213+
name := "kafka-connect-elastic-common",
211214
description := "Kafka Connect compatible connectors to move data between Kafka and popular data stores",
212-
libraryDependencies ++= baseDeps ++ kafkaConnectElastic6Deps,
215+
libraryDependencies ++= baseDeps ++ kafkaConnectElasticBaseDeps,
213216
publish / skip := true,
214-
FunctionalTest / baseDirectory := (LocalRootProject / baseDirectory).value,
215217
packExcludeJars := Seq(
216218
"scala-.*\\.jar",
217219
"zookeeper-.*\\.jar",
@@ -220,19 +222,20 @@ lazy val elastic6 = (project in file("kafka-connect-elastic6"))
220222
)
221223
.configureAssembly(true)
222224
.configureTests(baseTestDeps)
223-
.configureIntegrationTests(kafkaConnectElastic6TestDeps)
225+
.configureIntegrationTests(kafkaConnectElastic8TestDeps)
224226
.configureFunctionalTests()
225-
.enablePlugins(PackPlugin)
227+
.disablePlugins(PackPlugin)
226228

227-
lazy val elastic7 = (project in file("kafka-connect-elastic7"))
229+
lazy val elastic8 = (project in file("kafka-connect-elastic8"))
228230
.dependsOn(common)
229-
.dependsOn(`test-common` % "fun->compile")
231+
.dependsOn(`elastic-common`)
232+
.dependsOn(`test-common` % "fun->compile;it->compile")
230233
.settings(
231234
settings ++
232235
Seq(
233-
name := "kafka-connect-elastic7",
236+
name := "kafka-connect-elastic8",
234237
description := "Kafka Connect compatible connectors to move data between Kafka and popular data stores",
235-
libraryDependencies ++= baseDeps ++ kafkaConnectElastic7Deps,
238+
libraryDependencies ++= baseDeps ++ kafkaConnectElastic8Deps,
236239
publish / skip := true,
237240
packExcludeJars := Seq(
238241
"scala-.*\\.jar",
@@ -242,10 +245,33 @@ lazy val elastic7 = (project in file("kafka-connect-elastic7"))
242245
)
243246
.configureAssembly(true)
244247
.configureTests(baseTestDeps)
245-
.configureIntegrationTests(kafkaConnectElastic7TestDeps)
248+
.configureIntegrationTests(kafkaConnectElastic8TestDeps)
246249
.configureFunctionalTests()
247250
.enablePlugins(PackPlugin)
248251

252+
lazy val opensearch = (project in file("kafka-connect-opensearch"))
253+
.dependsOn(common)
254+
.dependsOn(`elastic-common`)
255+
.dependsOn(`test-common` % "fun->compile;it->compile")
256+
.settings(
257+
settings ++
258+
Seq(
259+
name := "kafka-connect-opensearch",
260+
description := "Kafka Connect compatible connectors to move data between Kafka and popular data stores",
261+
libraryDependencies ++= baseDeps ++ kafkaConnectOpenSearchDeps,
262+
publish / skip := true,
263+
packExcludeJars := Seq(
264+
"scala-.*\\.jar",
265+
"zookeeper-.*\\.jar",
266+
),
267+
),
268+
)
269+
.configureAssembly(false)
270+
.configureTests(baseTestDeps)
271+
//.configureIntegrationTests(kafkaConnectOpenSearchTestDeps)
272+
.configureFunctionalTests(bouncyCastle)
273+
.enablePlugins(PackPlugin)
274+
249275
lazy val http = (project in file("kafka-connect-http"))
250276
.dependsOn(common)
251277
//.dependsOn(`test-common` % "fun->compile")

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/sink/config/S3SinkConfig.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ object S3SinkConfig {
4949
s3ConfigDefBuilder.getInt(SEEK_MAX_INDEX_FILES),
5050
)
5151
} yield S3SinkConfig(
52-
S3Config(s3ConfigDefBuilder.getParsedValues),
52+
S3Config(s3ConfigDefBuilder.props),
5353
sinkBucketOptions,
5454
offsetSeekerOptions,
5555
s3ConfigDefBuilder.getCompressionCodec(),

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/sink/config/S3SinkConfigDefBuilder.scala

Lines changed: 3 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -18,19 +18,15 @@ package io.lenses.streamreactor.connect.aws.s3.sink.config
1818
import io.lenses.streamreactor.common.config.base.traits._
1919
import io.lenses.streamreactor.connect.aws.s3.config.DeleteModeSettings
2020
import io.lenses.streamreactor.connect.aws.s3.config.S3ConfigSettings
21+
import io.lenses.streamreactor.connect.cloud.common.config.CompressionCodecSettings
2122
import io.lenses.streamreactor.connect.cloud.common.sink.config.CloudSinkConfigDefBuilder
2223

23-
import scala.jdk.CollectionConverters.MapHasAsScala
24-
2524
case class S3SinkConfigDefBuilder(props: Map[String, String])
2625
extends BaseConfig(S3ConfigSettings.CONNECTOR_PREFIX, S3SinkConfigDef.config, props)
2726
with CloudSinkConfigDefBuilder
2827
with ErrorPolicySettings
2928
with NumberRetriesSettings
3029
with UserSettings
3130
with ConnectionSettings
32-
with DeleteModeSettings {
33-
34-
def getParsedValues: Map[String, _] = values().asScala.toMap
35-
36-
}
31+
with CompressionCodecSettings
32+
with DeleteModeSettings {}

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/source/S3SourceTask.scala

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -63,9 +63,8 @@ class S3SourceTask extends SourceTask with LazyLogging {
6363

6464
logger.debug(s"Received call to S3SourceTask.start with ${props.size()} properties")
6565

66-
val contextProperties: Map[String, String] =
67-
Option(context).flatMap(c => Option(c.configs()).map(_.asScala.toMap)).getOrElse(Map.empty)
68-
val mergedProperties: Map[String, String] = MapUtils.mergeProps(contextProperties, props.asScala.toMap)
66+
val contextProperties = Option(context).flatMap(c => Option(c.configs()).map(_.asScala.toMap)).getOrElse(Map.empty)
67+
val mergedProperties = MapUtils.mergeProps(contextProperties, props.asScala.toMap)
6968
(for {
7069
result <- S3SourceState.make(mergedProperties, contextOffsetFn)
7170
fiber <- result.partitionDiscoveryLoop.start

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/source/config/S3SourceConfig.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ object S3SourceConfig {
5050
S3SourceConfig(S3SourceConfigDefBuilder(props))
5151

5252
def apply(s3ConfigDefBuilder: S3SourceConfigDefBuilder): Either[Throwable, S3SourceConfig] = {
53-
val parsedValues = s3ConfigDefBuilder.getParsedValues
53+
val parsedValues = s3ConfigDefBuilder.props
5454
for {
5555
sbo <- SourceBucketOptions(
5656
s3ConfigDefBuilder,

kafka-connect-aws-s3/src/main/scala/io/lenses/streamreactor/connect/aws/s3/source/config/S3SourceConfigDefBuilder.scala

Lines changed: 1 addition & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -20,8 +20,6 @@ import io.lenses.streamreactor.connect.aws.s3.config.DeleteModeSettings
2020
import io.lenses.streamreactor.connect.aws.s3.config.S3ConfigSettings
2121
import io.lenses.streamreactor.connect.cloud.common.config.CompressionCodecSettings
2222

23-
import scala.jdk.CollectionConverters.MapHasAsScala
24-
2523
case class S3SourceConfigDefBuilder(props: Map[String, String])
2624
extends BaseConfig(S3ConfigSettings.CONNECTOR_PREFIX, S3SourceConfigDef.config, props)
2725
with KcqlSettings
@@ -31,8 +29,4 @@ case class S3SourceConfigDefBuilder(props: Map[String, String])
3129
with ConnectionSettings
3230
with CompressionCodecSettings
3331
with SourcePartitionSearcherSettings
34-
with DeleteModeSettings {
35-
36-
def getParsedValues: Map[String, _] = values().asScala.toMap
37-
38-
}
32+
with DeleteModeSettings {}

kafka-connect-azure-datalake/src/main/scala/io/lenses/streamreactor/connect/datalake/sink/config/DatalakeSinkConfig.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -49,7 +49,7 @@ object DatalakeSinkConfig {
4949
s3ConfigDefBuilder.getInt(SEEK_MAX_INDEX_FILES),
5050
)
5151
} yield DatalakeSinkConfig(
52-
AzureConfig(s3ConfigDefBuilder.getParsedValues, authMode),
52+
AzureConfig(s3ConfigDefBuilder.props, authMode),
5353
sinkBucketOptions,
5454
offsetSeekerOptions,
5555
s3ConfigDefBuilder.getCompressionCodec(),

kafka-connect-azure-datalake/src/main/scala/io/lenses/streamreactor/connect/datalake/sink/config/DatalakeSinkConfigDefBuilder.scala

Lines changed: 3 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -16,21 +16,17 @@
1616
package io.lenses.streamreactor.connect.datalake.sink.config
1717

1818
import io.lenses.streamreactor.common.config.base.traits._
19+
import io.lenses.streamreactor.connect.cloud.common.config.CompressionCodecSettings
1920
import io.lenses.streamreactor.connect.cloud.common.sink.config.CloudSinkConfigDefBuilder
2021
import io.lenses.streamreactor.connect.datalake.config.AuthModeSettings
2122
import io.lenses.streamreactor.connect.datalake.config.AzureConfigSettings
2223

23-
import scala.jdk.CollectionConverters.MapHasAsScala
24-
2524
case class DatalakeSinkConfigDefBuilder(props: Map[String, String])
2625
extends BaseConfig(AzureConfigSettings.CONNECTOR_PREFIX, DatalakeSinkConfigDef.config, props)
2726
with CloudSinkConfigDefBuilder
2827
with ErrorPolicySettings
2928
with NumberRetriesSettings
3029
with UserSettings
3130
with ConnectionSettings
32-
with AuthModeSettings {
33-
34-
def getParsedValues: Map[String, _] = values().asScala.toMap
35-
36-
}
31+
with CompressionCodecSettings
32+
with AuthModeSettings {}

kafka-connect-azure-documentdb/src/main/scala/io/lenses/streamreactor/connect/azure/documentdb/sink/DocumentDbSinkConnector.scala

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
*/
1616
package io.lenses.streamreactor.connect.azure.documentdb.sink
1717

18+
import cats.implicits.toBifunctorOps
1819
import io.lenses.streamreactor.common.config.Helpers
1920
import io.lenses.streamreactor.common.utils.JarManifest
2021
import io.lenses.streamreactor.connect.azure.documentdb.DocumentClientProvider
@@ -100,7 +101,7 @@ class DocumentDbSinkConnector private[sink] (builder: DocumentDbSinkSettings =>
100101
configProps = props
101102

102103
//check input topics
103-
Helpers.checkInputTopics(DocumentDbConfigConstants.KCQL_CONFIG, props.asScala.toMap)
104+
Helpers.checkInputTopics(DocumentDbConfigConstants.KCQL_CONFIG, props.asScala.toMap).leftMap(throw _)
104105

105106
val settings = DocumentDbSinkSettings(config)
106107

kafka-connect-cassandra/src/main/scala/io/lenses/streamreactor/connect/cassandra/CassandraConnection.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -15,9 +15,9 @@
1515
*/
1616
package io.lenses.streamreactor.connect.cassandra
1717

18-
import io.lenses.streamreactor.common.config.SSLConfig
19-
import io.lenses.streamreactor.common.config.SSLConfigContext
2018
import io.lenses.streamreactor.connect.cassandra.config.CassandraConfigConstants
19+
import io.lenses.streamreactor.connect.cassandra.config.SSLConfig
20+
import io.lenses.streamreactor.connect.cassandra.config.SSLConfigContext
2121
import io.lenses.streamreactor.connect.cassandra.config.LoadBalancingPolicy
2222
import com.datastax.driver.core.Cluster.Builder
2323
import com.datastax.driver.core.policies.DCAwareRoundRobinPolicy

0 commit comments

Comments
 (0)