diff --git a/airbyte-cdk/java/airbyte-cdk/datastore-mongo/src/main/kotlin/io/airbyte/cdk/db/mongodb/MongoDatabase.kt b/airbyte-cdk/java/airbyte-cdk/datastore-mongo/src/main/kotlin/io/airbyte/cdk/db/mongodb/MongoDatabase.kt index 7b36576ff5bd..02f11b3d0802 100644 --- a/airbyte-cdk/java/airbyte-cdk/datastore-mongo/src/main/kotlin/io/airbyte/cdk/db/mongodb/MongoDatabase.kt +++ b/airbyte-cdk/java/airbyte-cdk/datastore-mongo/src/main/kotlin/io/airbyte/cdk/db/mongodb/MongoDatabase.kt @@ -16,7 +16,6 @@ import io.airbyte.commons.util.MoreIterators import java.util.* import java.util.Spliterators.AbstractSpliterator import java.util.function.Consumer -import java.util.stream.Collectors import java.util.stream.Stream import java.util.stream.StreamSupport import org.bson.BsonDocument @@ -57,9 +56,8 @@ class MongoDatabase(connectionString: String, databaseName: String) : get() { val collectionNames = database.listCollectionNames() ?: return Collections.emptySet() return MoreIterators.toSet(collectionNames.iterator()) - .stream() .filter { c: String -> !c.startsWith(MONGO_RESERVED_COLLECTION_PREFIX) } - .collect(Collectors.toSet()) + .toSet() } fun getCollection(collectionName: String): MongoCollection { diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/jdbc/AbstractJdbcSource.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/jdbc/AbstractJdbcSource.kt index 40991c414416..03feabc69b86 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/jdbc/AbstractJdbcSource.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/jdbc/AbstractJdbcSource.kt @@ -365,10 +365,7 @@ abstract class AbstractJdbcSource( ): Map> { LOGGER.info( "Discover primary keys for tables: " + - tableInfos - .stream() - .map { obj: TableInfo> -> obj.name } - .collect(Collectors.toSet()) + tableInfos.map { obj: TableInfo> -> obj.name }.toSet() ) try { // Get all primary keys without specifying a table name diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/DbSourceDiscoverUtil.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/DbSourceDiscoverUtil.kt index 40c99432deef..7b6a023062a2 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/DbSourceDiscoverUtil.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/DbSourceDiscoverUtil.kt @@ -143,9 +143,8 @@ object DbSourceDiscoverUtil { ) .withSourceDefinedPrimaryKey(primaryKeys) } - // This is ugly. Some of our tests change the streams on the AirbyteCatalog - // object... - .toMutableList() + .toMutableList() // This is ugly, but we modify this list in + // JdbcSourceAcceptanceTest.testDiscoverWithMultipleSchemas return AirbyteCatalog().withStreams(streams) } diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/RelationalDbReadUtil.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/RelationalDbReadUtil.kt index 2ea21ddcc7ac..255ca51c885e 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/RelationalDbReadUtil.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/RelationalDbReadUtil.kt @@ -9,7 +9,6 @@ import io.airbyte.protocol.models.v0.AirbyteStreamNameNamespacePair import io.airbyte.protocol.models.v0.ConfiguredAirbyteCatalog import io.airbyte.protocol.models.v0.ConfiguredAirbyteStream import io.airbyte.protocol.models.v0.SyncMode -import java.util.stream.Collectors object RelationalDbReadUtil { fun identifyStreamsToSnapshot( @@ -38,11 +37,10 @@ object RelationalDbReadUtil { ): List { val initialLoadStreamsNamespacePairs = streamsForInitialLoad - .stream() .map { stream: ConfiguredAirbyteStream -> AirbyteStreamNameNamespacePair.fromAirbyteStream(stream.stream) } - .collect(Collectors.toSet()) + .toSet() return catalog.streams .stream() .filter { c: ConfiguredAirbyteStream -> c.syncMode == SyncMode.INCREMENTAL } diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/CursorManager.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/CursorManager.kt index 90a086a68a14..3f78b9726e72 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/CursorManager.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/CursorManager.kt @@ -92,7 +92,6 @@ class CursorManager( ): Map { val allStreamNames = catalog.streams - .stream() .filter { c: ConfiguredAirbyteStream -> if (onlyIncludeIncrementalStreams) { return@filter c.syncMode == SyncMode.INCREMENTAL @@ -103,7 +102,7 @@ class CursorManager( .map { stream: AirbyteStream -> AirbyteStreamNameNamespacePair.fromAirbyteStream(stream) } - .collect(Collectors.toSet()) + .toMutableSet() allStreamNames.addAll( streamSupplier .get() diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/GlobalStateManager.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/GlobalStateManager.kt index a4c475b4bc06..9948a578b5dd 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/GlobalStateManager.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/main/kotlin/io/airbyte/cdk/integrations/source/relationaldb/state/GlobalStateManager.kt @@ -105,7 +105,6 @@ class GlobalStateManager( ): Set { if (airbyteStateMessage!!.type == AirbyteStateMessage.AirbyteStateType.GLOBAL) { return airbyteStateMessage.global.streamStates - .stream() .map { streamState: AirbyteStreamState -> val cloned = Jsons.clone(streamState) AirbyteStreamNameNamespacePair( @@ -113,7 +112,7 @@ class GlobalStateManager( cloned.streamDescriptor.namespace ) } - .collect(Collectors.toSet()) + .toSet() } else { val legacyState: DbState? = Jsons.`object`(airbyteStateMessage.data, DbState::class.java) @@ -127,12 +126,11 @@ class GlobalStateManager( streams: List ): Set { return streams - .stream() .map { stream: DbStreamState -> val cloned = Jsons.clone(stream) AirbyteStreamNameNamespacePair(cloned.streamName, cloned.streamNamespace) } - .collect(Collectors.toSet()) + .toSet() } companion object { diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/debezium/CdcSourceTest.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/debezium/CdcSourceTest.kt index 996a369794a1..ca5e209b9579 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/debezium/CdcSourceTest.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/debezium/CdcSourceTest.kt @@ -15,8 +15,6 @@ import io.airbyte.protocol.models.JsonSchemaType import io.airbyte.protocol.models.v0.* import java.util.* import java.util.function.Consumer -import java.util.stream.Collectors -import java.util.stream.Stream import org.junit.jupiter.api.AfterEach import org.junit.jupiter.api.Assertions import org.junit.jupiter.api.BeforeEach @@ -267,7 +265,7 @@ abstract class CdcSourceTest> { if (message.type == AirbyteMessage.Type.RECORD) { val recordMessage = message.record recordsPerStream - .computeIfAbsent(recordMessage.stream) { c: String -> ArrayList() } + .computeIfAbsent(recordMessage.stream) { _: String -> ArrayList() } .add(recordMessage) } } @@ -303,10 +301,7 @@ abstract class CdcSourceTest> { assertExpectedRecords( expectedRecords, actualRecords, - actualRecords - .stream() - .map { obj: AirbyteRecordMessage -> obj.stream } - .collect(Collectors.toSet()) + actualRecords.map { obj: AirbyteRecordMessage -> obj.stream }.toSet() ) } @@ -333,8 +328,7 @@ abstract class CdcSourceTest> { ) { val actualData = actualRecords - .stream() - .map { recordMessage: AirbyteRecordMessage -> + .map { recordMessage: AirbyteRecordMessage -> Assertions.assertTrue(streamNames.contains(recordMessage.stream)) Assertions.assertNotNull(recordMessage.emittedAt) @@ -351,7 +345,7 @@ abstract class CdcSourceTest> { removeCDCColumns(data as ObjectNode) data } - .collect(Collectors.toSet()) + .toSet() Assertions.assertEquals(expectedRecords, actualData) } @@ -604,8 +598,7 @@ abstract class CdcSourceTest> { // Full refresh does not get any state messages. assertExpectedStateMessageCountMatches(stateMessages1, MODEL_RECORDS_2.size.toLong()) assertExpectedRecords( - Streams.concat(MODEL_RECORDS_2.stream(), MODEL_RECORDS.stream()) - .collect(Collectors.toSet()), + (MODEL_RECORDS_2 + MODEL_RECORDS).toSet(), recordMessages1, setOf(MODELS_STREAM_NAME), names, @@ -626,8 +619,7 @@ abstract class CdcSourceTest> { assertExpectedStateMessagesFromIncrementalSync(stateMessages2) assertExpectedStateMessageCountMatches(stateMessages2, 1) assertExpectedRecords( - Streams.concat(MODEL_RECORDS_2.stream(), Stream.of(puntoRecord)) - .collect(Collectors.toSet()), + (MODEL_RECORDS_2 + puntoRecord).toSet(), recordMessages2, setOf(MODELS_STREAM_NAME), names, @@ -721,9 +713,8 @@ abstract class CdcSourceTest> { Assertions.assertNotNull(stateMessageEmittedAfterFirstSyncCompletion.global.sharedState) val streamsInStateAfterFirstSyncCompletion = stateMessageEmittedAfterFirstSyncCompletion.global.streamStates - .stream() .map { obj: AirbyteStreamState -> obj.streamDescriptor } - .collect(Collectors.toSet()) + .toSet() Assertions.assertEquals(1, streamsInStateAfterFirstSyncCompletion.size) Assertions.assertTrue( streamsInStateAfterFirstSyncCompletion.contains( @@ -813,9 +804,8 @@ abstract class CdcSourceTest> { HashSet(MODEL_RECORDS_RANDOM), recordsForModelsRandomStreamFromSecondBatch, recordsForModelsRandomStreamFromSecondBatch - .stream() .map { obj: AirbyteRecordMessage -> obj.stream } - .collect(Collectors.toSet()), + .toSet(), Sets.newHashSet(RANDOM_TABLE_NAME), randomSchema() ) @@ -882,9 +872,8 @@ abstract class CdcSourceTest> { ) val streamsInSyncCompletionStateAfterThirdSync = stateMessageEmittedAfterThirdSyncCompletion.global.streamStates - .stream() .map { obj: AirbyteStreamState -> obj.streamDescriptor } - .collect(Collectors.toSet()) + .toSet() Assertions.assertTrue( streamsInSyncCompletionStateAfterThirdSync.contains( StreamDescriptor().withName(RANDOM_TABLE_NAME).withNamespace(randomSchema()) @@ -913,9 +902,8 @@ abstract class CdcSourceTest> { recordsWrittenInRandomTable, recordsForModelsRandomStreamFromThirdBatch, recordsForModelsRandomStreamFromThirdBatch - .stream() .map { obj: AirbyteRecordMessage -> obj.stream } - .collect(Collectors.toSet()), + .toSet(), Sets.newHashSet(RANDOM_TABLE_NAME), randomSchema() ) @@ -937,9 +925,8 @@ abstract class CdcSourceTest> { ) val streamsInSnapshotState = stateMessageEmittedAfterSnapshotCompletionInSecondSync.global.streamStates - .stream() .map { obj: AirbyteStreamState -> obj.streamDescriptor } - .collect(Collectors.toSet()) + .toSet() Assertions.assertEquals(2, streamsInSnapshotState.size) Assertions.assertTrue( streamsInSnapshotState.contains( @@ -964,9 +951,8 @@ abstract class CdcSourceTest> { ) val streamsInSyncCompletionState = stateMessageEmittedAfterSecondSyncCompletion.global.streamStates - .stream() .map { obj: AirbyteStreamState -> obj.streamDescriptor } - .collect(Collectors.toSet()) + .toSet() Assertions.assertEquals(2, streamsInSnapshotState.size) Assertions.assertTrue( streamsInSyncCompletionState.contains( diff --git a/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/source/jdbc/test/JdbcStressTest.kt b/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/source/jdbc/test/JdbcStressTest.kt index f231d01ab0c2..9fbd2d1252fc 100644 --- a/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/source/jdbc/test/JdbcStressTest.kt +++ b/airbyte-cdk/java/airbyte-cdk/db-sources/src/testFixtures/kotlin/io/airbyte/cdk/integrations/source/jdbc/test/JdbcStressTest.kt @@ -20,7 +20,6 @@ import io.airbyte.protocol.models.Field import io.airbyte.protocol.models.JsonSchemaType import io.airbyte.protocol.models.v0.* import java.math.BigDecimal -import java.nio.ByteBuffer import java.sql.Connection import java.util.* import org.junit.jupiter.api.Assertions @@ -172,7 +171,6 @@ abstract class JdbcStressTest { } .peek { m: AirbyteMessage -> assertExpectedMessage(m) } .count() - var a: ByteBuffer val expectedRoundedRecordsCount = TOTAL_RECORDS - TOTAL_RECORDS % 1000 LOGGER.info("expected records count: " + TOTAL_RECORDS) LOGGER.info("actual records count: $actualCount") diff --git a/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/commons/enums/Enums.kt b/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/commons/enums/Enums.kt index 56a964761d41..9bee126061c8 100644 --- a/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/commons/enums/Enums.kt +++ b/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/commons/enums/Enums.kt @@ -7,11 +7,9 @@ package io.airbyte.commons.enums import com.google.common.base.Preconditions import com.google.common.collect.Maps import com.google.common.collect.Sets -import java.util.Arrays import java.util.Locale import java.util.Optional import java.util.concurrent.ConcurrentMap -import java.util.stream.Collectors class Enums { companion object { @@ -54,12 +52,8 @@ class Enums { Preconditions.checkArgument(c2.isEnum) return (c1.enumConstants.size == c2.enumConstants.size && Sets.difference( - Arrays.stream(c1.enumConstants) - .map { obj: T1 -> obj!!.name } - .collect(Collectors.toSet()), - Arrays.stream(c2.enumConstants) - .map { obj: T2 -> obj!!.name } - .collect(Collectors.toSet()), + c1.enumConstants.map { obj: T1 -> obj!!.name }.toSet(), + c2.enumConstants.map { obj: T2 -> obj!!.name }.toSet(), ) .isEmpty()) } diff --git a/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/validation/json/JsonSchemaValidator.kt b/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/validation/json/JsonSchemaValidator.kt index a7b85b909708..c06001032560 100644 --- a/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/validation/json/JsonSchemaValidator.kt +++ b/airbyte-cdk/java/airbyte-cdk/dependencies/src/main/kotlin/io/airbyte/validation/json/JsonSchemaValidator.kt @@ -11,7 +11,6 @@ import java.io.File import java.io.IOException import java.net.URI import java.net.URISyntaxException -import java.util.stream.Collectors import me.andrz.jackson.JsonContext import me.andrz.jackson.JsonReferenceException import me.andrz.jackson.JsonReferenceProcessor @@ -89,9 +88,8 @@ class JsonSchemaValidator @VisibleForTesting constructor(private val baseUri: UR fun validate(schemaJson: JsonNode, objectJson: JsonNode): Set { return validateInternal(schemaJson, objectJson) - .stream() .map { obj: ValidationMessage -> obj.message } - .collect(Collectors.toSet()) + .toSet() } fun getValidationMessageArgs(schemaJson: JsonNode, objectJson: JsonNode): List> { diff --git a/airbyte-cdk/java/airbyte-cdk/dependencies/src/testFixtures/kotlin/io/airbyte/workers/helper/CatalogClientConverters.kt b/airbyte-cdk/java/airbyte-cdk/dependencies/src/testFixtures/kotlin/io/airbyte/workers/helper/CatalogClientConverters.kt index 619f5be63677..a480ef7bab05 100644 --- a/airbyte-cdk/java/airbyte-cdk/dependencies/src/testFixtures/kotlin/io/airbyte/workers/helper/CatalogClientConverters.kt +++ b/airbyte-cdk/java/airbyte-cdk/dependencies/src/testFixtures/kotlin/io/airbyte/workers/helper/CatalogClientConverters.kt @@ -10,7 +10,6 @@ import io.airbyte.commons.text.Names import io.airbyte.protocol.models.SyncMode import io.airbyte.validation.json.JsonValidationException import java.util.* -import java.util.stream.Collectors /** * Utilities to convert Catalog protocol to Catalog API client. This class was similar to existing @@ -76,9 +75,8 @@ object CatalogClientConverters { // field path. val selectedFieldNames = config.selectedFields!! - .stream() .map { field: SelectedFieldInfo -> field.fieldPath!![0] } - .collect(Collectors.toSet()) + .toSet() // TODO(mfsiega-airbyte): we only check the top level of the cursor/primary key fields // because we // don't support filtering nested fields yet. diff --git a/airbyte-cdk/java/airbyte-cdk/gcs-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/gcs/GcsAvroParquetDestinationAcceptanceTest.kt b/airbyte-cdk/java/airbyte-cdk/gcs-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/gcs/GcsAvroParquetDestinationAcceptanceTest.kt index aa1ec7c96221..05fb1799c1dc 100644 --- a/airbyte-cdk/java/airbyte-cdk/gcs-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/gcs/GcsAvroParquetDestinationAcceptanceTest.kt +++ b/airbyte-cdk/java/airbyte-cdk/gcs-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/gcs/GcsAvroParquetDestinationAcceptanceTest.kt @@ -148,10 +148,9 @@ abstract class GcsAvroParquetDestinationAcceptanceTest(fileUploadFormat: FileUpl field .schema() .types - .stream() .map { obj: Schema -> obj.type } .filter { type: Schema.Type -> type != Schema.Type.NULL } - .collect(Collectors.toSet()) + .toSet() } ) ) @@ -165,12 +164,11 @@ abstract class GcsAvroParquetDestinationAcceptanceTest(fileUploadFormat: FileUpl field .schema() .types - .stream() .filter { type: Schema -> type.type != Schema.Type.NULL } - .flatMap { type: Schema -> type.elementType.types.stream() } + .flatMap { type: Schema -> type.elementType.types } .map { obj: Schema -> obj.type } .filter { type: Schema.Type -> type != Schema.Type.NULL } - .collect(Collectors.toSet()) + .toSet() } ) ) diff --git a/airbyte-cdk/java/airbyte-cdk/s3-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/s3/S3AvroParquetDestinationAcceptanceTest.kt b/airbyte-cdk/java/airbyte-cdk/s3-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/s3/S3AvroParquetDestinationAcceptanceTest.kt index 19547cce3b74..8e0109311e7f 100644 --- a/airbyte-cdk/java/airbyte-cdk/s3-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/s3/S3AvroParquetDestinationAcceptanceTest.kt +++ b/airbyte-cdk/java/airbyte-cdk/s3-destinations/src/testFixtures/kotlin/io/airbyte/cdk/integrations/destination/s3/S3AvroParquetDestinationAcceptanceTest.kt @@ -79,13 +79,14 @@ protected constructor(fileUploadFormat: FileUploadFormat) : else fieldDefinition["type"] val airbyteTypeProperty = fieldDefinition["airbyte_type"] val airbyteTypePropertyText = airbyteTypeProperty?.asText() - return Arrays.stream(JsonSchemaType.entries.toTypedArray()) + return JsonSchemaType.entries + .toTypedArray() .filter { value: JsonSchemaType -> value.jsonSchemaType == typeProperty.asText() && compareAirbyteTypes(airbyteTypePropertyText, value) } .map { obj: JsonSchemaType -> obj.avroType } - .collect(Collectors.toSet()) + .toSet() } private fun compareAirbyteTypes( @@ -138,10 +139,9 @@ protected constructor(fileUploadFormat: FileUploadFormat) : field .schema() .types - .stream() .map { obj: Schema -> obj.type } .filter { type: Schema.Type -> type != Schema.Type.NULL } - .collect(Collectors.toSet()) + .toSet() } ) ) @@ -155,12 +155,11 @@ protected constructor(fileUploadFormat: FileUploadFormat) : field .schema() .types - .stream() .filter { type: Schema -> type.type != Schema.Type.NULL } - .flatMap { type: Schema -> type.elementType.types.stream() } + .flatMap { type: Schema -> type.elementType.types } .map { obj: Schema -> obj.type } .filter { type: Schema.Type -> type != Schema.Type.NULL } - .collect(Collectors.toSet()) + .toSet() } ) )