diff --git a/.github/workflows/ds-release.yml b/.github/workflows/ds-release.yml new file mode 100644 index 00000000..26c3f37d --- /dev/null +++ b/.github/workflows/ds-release.yml @@ -0,0 +1,26 @@ +name: DS Release + +on: + push: + tags: + - 'v*' + +jobs: + create-release: + permissions: write-all + runs-on: ubuntu-latest + + steps: + - uses: actions/checkout@v2 + - name: Set up JDK 17 + uses: actions/setup-java@v2 + with: + java-version: '17' + distribution: 'adopt' + - name: Build with Maven + run: mvn -B package -DskipTests -Dcheckstyle.skip + - uses: ncipollo/release-action@v1 + with: + artifacts: "**/target/*.nar" + token: ${{ secrets.GITHUB_TOKEN }} + generateReleaseNotes: true \ No newline at end of file diff --git a/.github/workflows/pr-test.yml b/.github/workflows/pr-test.yml index 2299ef85..fc5fdd15 100644 --- a/.github/workflows/pr-test.yml +++ b/.github/workflows/pr-test.yml @@ -4,7 +4,7 @@ on: pull_request: push: branches: - - master + - 3.2_ds jobs: build: diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index d1d54dfc..bb053467 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -36,7 +36,7 @@ jobs: PULSAR_VERSION=`mvn -q -Dexec.executable=echo -Dexec.args='${pulsar.version}' --non-recursive exec:exec 2>/dev/null` REPO=`mvn -q -Dexec.executable=echo -Dexec.args='${project.artifactId}' --non-recursive exec:exec 2>/dev/null` IMAGE_REPO=streamnative/${REPO} - RUNNER_IMAGE=docker.cloudsmith.io/streamnative/staging/pulsar-functions-java-runner:${PULSAR_VERSION} + RUNNER_IMAGE=streamnative/pulsar-functions-java-runner:${PULSAR_VERSION} docker pull ${RUNNER_IMAGE} docker build --build-arg PULSAR_VERSION="$PULSAR_VERSION" -t ${IMAGE_REPO}:${CONNECTOR_VERSION} -f ./image/Dockerfile ./ docker push ${IMAGE_REPO}:${CONNECTOR_VERSION} diff --git a/.github/workflows/sonarqube.yaml b/.github/workflows/sonarqube.yaml new file mode 100644 index 00000000..35506dfa --- /dev/null +++ b/.github/workflows/sonarqube.yaml @@ -0,0 +1,106 @@ +name: SonarQube Scan + +on: + push: + branches: + - 4.0_ds + - 3.2_ds + pull_request: + types: [opened, synchronize, reopened] + branches: + - 4.0_ds + - 3.2_ds + workflow_dispatch: + inputs: + scan_type: + description: 'Type of scan to run' + required: true + type: choice + options: + - branch + - pr + default: branch + ref: + description: 'Commit SHA, branch, or tag to check out and scan (leave blank to use the default branch)' + required: false + type: string + default: '' + +concurrency: + group: sonarqube-${{ github.ref }} + cancel-in-progress: true + +env: + TRUSTSTORE_PATH: ./ibm_castorevpcprod + # Maven tuning — mirrors settings used across the rest of the CI workflows. + MAVEN_OPTS: >- + -Xss1500k + -Xmx2048m + -XX:+UnlockDiagnosticVMOptions + -XX:+IgnoreUnrecognizedVMOptions + -XX:GCLockerRetryAllocationCount=100 + -Daether.connector.http.reuseConnections=false + -Daether.connector.requestTimeout=60000 + -Dhttp.keepAlive=false + -Dmaven.wagon.http.pool=false + -Dmaven.wagon.http.retryHandler.class=standard + -Dmaven.wagon.http.retryHandler.count=3 + -Dmaven.wagon.http.retryHandler.requestSentEnabled=true + -Dmaven.wagon.http.serviceUnavailableRetryStrategy.class=standard + -Dmaven.wagon.rto=60000 + +jobs: + sonarqube-scan: + name: SonarQube Scan + runs-on: ubuntu-latest + + steps: + - name: Checkout repository + uses: actions/checkout@v4 + with: + # Full history required for accurate blame and new-code period detection. + fetch-depth: 0 + # When triggered manually with a specific ref (commit SHA, branch, tag), + # check out that ref; otherwise fall back to the default event ref. + ref: ${{ (github.event_name == 'workflow_dispatch' && inputs.ref != '') && inputs.ref || github.ref }} + + - name: Set up JDK 17 + uses: actions/setup-java@v4 + with: + distribution: corretto + java-version: '17' + cache: maven + + # Decode the IBM truststore from the base64 secret and write it to disk. + # Secret is bound to an env var to avoid shell-injection (S7636). + - name: Decode IBM SonarQube truststore + env: + TRUSTSTORE_B64: ${{ secrets.IBM_SONARQUBE_CERT_TRUSTSTORE_BASE64 }} + run: | + echo "$TRUSTSTORE_B64" | base64 --decode > "${{ env.TRUSTSTORE_PATH }}" + + # Produce .class files required by SonarQube's Java bytecode analyser. + - name: Build + run: mvn -B package -DskipTests --no-transfer-progress + + # Pinned to v8.2.1 (SHA: 22918119ff8e1ca75a623e15c8296b6ea4fbe28f). + - name: Run SonarQube scan + uses: SonarSource/sonarqube-scan-action@22918119ff8e1ca75a623e15c8296b6ea4fbe28f + env: + SONAR_TOKEN: ${{ secrets.IBM_SONARQUBE_API_TOKEN }} + SONAR_HOST_URL: ${{ vars.SONAR_HOST_URL }} + with: + args: > + -D sonar.projectKey=${{ vars.SONAR_PROJECT_KEY }} + -D sonar.projectName=${{ vars.SONAR_PROJECT_NAME }} + -D sonar.java.source=17 + -D sonar.sources=. + -D sonar.exclusions=**/proto/**,**/shade/**,**/shaded/**,**/target/** + -D sonar.java.binaries=**/target/classes + -D sonar.scanner.truststorePath=${{ env.TRUSTSTORE_PATH }} + -D sonar.scanner.truststorePassword=${{ secrets.IBM_SONARQUBE_CERT_PASSWORD }} + ${{ (github.event_name == 'pull_request' || (github.event_name == 'workflow_dispatch' && inputs.scan_type == 'pr')) + && format('-Dsonar.pullrequest.key={0} -Dsonar.pullrequest.branch={1} -Dsonar.pullrequest.base={2}', + github.event.pull_request.number, github.head_ref, github.base_ref) + || format('-Dsonar.branch.name={0}', + (github.event_name == 'workflow_dispatch' && inputs.ref != '') && inputs.ref || github.ref_name) }} \ No newline at end of file diff --git a/image/Dockerfile b/image/Dockerfile index 0c0e2ea7..aa9ce43e 100644 --- a/image/Dockerfile +++ b/image/Dockerfile @@ -18,5 +18,5 @@ # ARG PULSAR_VERSION -FROM docker.cloudsmith.io/streamnative/staging/pulsar-functions-java-runner:${PULSAR_VERSION} +FROM streamnative/pulsar-functions-java-runner:${PULSAR_VERSION} COPY --chown=$UID:$GID target/*.nar /pulsar/connectors/ diff --git a/pom.xml b/pom.xml index 6afb0501..91b1176d 100644 --- a/pom.xml +++ b/pom.xml @@ -18,8 +18,7 @@ under the License. --> - + org.apache apache @@ -30,7 +29,7 @@ io.streamnative.ecosystem pulsar-io-cloud-storage - 2.9.0-rc-202110221101 + 3.2.11-SNAPSHOT Pulsar Ecosystem :: IO Connector :: Cloud Storage Project Cloud Storage Connector integrates Apache Pulsar with Cloud Storage. @@ -45,23 +44,37 @@ 2 - 2.13.2 - 2.13.4.1 - 1.18.20 - 2.11.0.4 - 1.12.2 - 1.12.148 - 2.16.1 + 2.18.6 + 2.21.2 + 1.18.32 + com.datastax.oss + 3.1.4.28 + 1.11.4 + 3.4.1 + 1.15.2 + 1.12.698 + 2.16.104 2.3.1 2.8.9 - 1.21 + 1.26.0 2.17.1 - 2.4.7 - 3.19.6 + 2.4.9 + 3.25.5 + 2.15.0 + 2.46 + 3.6.0 + 1.5.4 + 1.1.10.4 + 9.37.4 + 1.11.0 + 2.15.0 + 2.14.0 + 2.0.3 + 1.84 ${protobuf3.version} 1.42.1 ${grpc.version} - 1.2.6 + 1.2.22 4.13.1 @@ -69,6 +82,7 @@ 3.7.7 2.0.2 1.15.2 + 3.18.0 3.0 @@ -83,8 +97,10 @@ 2.0 1.2.24 5.4.0 - 4.1.86.Final - 0.17.0 + 4.1.135.Final + 2.1.0 + 0.23.0 + 5.3.5 @@ -95,19 +111,36 @@ + + scm:git:git@github.com:datastax/pulsar-io-cloud-storage.git + scm:git:git@github.com:datastax/pulsar-io-cloud-storage.git + https://github.com/datastax/pulsar-io-cloud-storage + HEAD + + - io.streamnative + ${pulsar.groupId} pulsar-io-core ${pulsar.version} - io.streamnative + ${pulsar.groupId} pulsar-io-common ${pulsar.version} + + ${pulsar.groupId} + pulsar-common + ${pulsar.version} + + + ${pulsar.groupId} + pulsar-client-original + ${pulsar.version} + com.fasterxml.jackson.core jackson-databind @@ -129,7 +162,7 @@ ${spotbugs-annotations.version} - io.streamnative + ${pulsar.groupId} jclouds-shaded ${pulsar.version} @@ -149,6 +182,99 @@ + + + org.asynchttpclient + async-http-client + ${asynchttpclient.version} + runtime + + + + org.glassfish.jersey.core + jersey-client + ${jersey-client.version} + runtime + + + + dnsjava + dnsjava + ${dnsjava.version} + runtime + + + + org.codehaus.jettison + jettison + ${jettison.version} + runtime + + + + org.xerial.snappy + snappy-java + ${snappy.version} + runtime + + + + com.nimbusds + nimbus-jose-jwt + ${nimbusds.version} + runtime + + + + commons-beanutils + commons-beanutils + ${commons-beanutils.version} + runtime + + + + org.apache.commons + commons-configuration2 + ${commons-configuration2.version} + runtime + + + + org.bouncycastle + bcpkix-jdk18on + ${bouncycastle.version} + runtime + + + + org.bouncycastle + bcprov-jdk18on + ${bouncycastle.version} + runtime + + + + org.jetbrains.kotlin + kotlin-stdlib + ${kotlin-stdlib.version} + runtime + + + + org.apache.httpcomponents.core5 + httpcore5-h2 + ${httpcore5-h2.version} + + + commons-io + commons-io + ${commons-io.version} + + + io.airlift + aircompressor + ${aircompressor.version} + com.google.protobuf protobuf-bom @@ -194,6 +320,11 @@ pom import + + org.apache.avro + avro + ${avro.version} + @@ -304,20 +435,34 @@ - io.streamnative + ${pulsar.groupId} pulsar-io-core + + ${pulsar.groupId} + pulsar-io-common + + + ${pulsar.groupId} + pulsar-common + + + ${pulsar.groupId} + pulsar-client-original + com.fasterxml.jackson.core jackson-databind + ${jackson-databind.version} com.fasterxml.jackson.dataformat jackson-dataformat-yaml + ${jackson.version} - io.streamnative + ${pulsar.groupId} jclouds-shaded @@ -385,39 +530,11 @@ org.apache.avro avro - 1.10.2 - - - org.apache.hadoop - hadoop-client - 3.3.3 - - - org.apache.hadoop - hadoop-yarn-api - - - org.apache.hadoop - hadoop-yarn-client - - - org.apache.hadoop - hadoop-yarn-common - - - org.apache.hadoop - hadoop-mapreduce-client-jobclient - - - org.apache.hadoop - hadoop-hdfs-client - - org.apache.hadoop hadoop-common - 3.3.3 + ${hadoop.version} org.apache.hadoop @@ -447,6 +564,41 @@ log4j log4j + + org.apache.zookeeper + zookeeper + + + + + org.apache.hadoop + hadoop-client + ${hadoop.version} + + + org.apache.hadoop + hadoop-yarn-api + + + org.apache.hadoop + hadoop-yarn-client + + + org.apache.hadoop + hadoop-yarn-common + + + org.apache.hadoop + hadoop-mapreduce-client-jobclient + + + org.apache.hadoop + hadoop-hdfs-client + + + org.apache.zookeeper + zookeeper + @@ -563,16 +715,10 @@ pulsar test - - org.apache.commons commons-lang3 - 3.11 - - - io.streamnative - pulsar-io-common + ${commons-lang3.version} javax.xml.bind @@ -580,12 +726,7 @@ ${jaxb-api} - io.streamnative - pulsar-client-original - ${pulsar.version} - - - io.streamnative + ${pulsar.groupId} pulsar-client-admin-original ${pulsar.version} @@ -595,11 +736,6 @@ 5.7.0 test - - io.streamnative - pulsar-common - ${pulsar.version} - @@ -611,6 +747,8 @@ ${maven-compiler-plugin.version} UTF-8 + ${java.version} + ${java.version} -Xlint:deprecation @@ -624,8 +762,9 @@ org.apache.maven.plugins maven-surefire-plugin + ${maven-surefire-plugin.version} - + false 1800 ${testRetryCount} @@ -821,11 +960,38 @@ central default https://repo1.maven.org/maven2 + + false + - bintray-streamnative-maven - bintray - https://dl.bintray.com/streamnative/maven + confluent + https://packages.confluent.io/maven/ + + false + + + + datastax-releases + https://repo.datastax.com/artifactory/datastax-public-releases-local + + false + + + true + + + + datastax-releases + DataStax Local Releases + https://repo.aws.dsinternal.org/artifactory/datastax-releases-local/ + + + datastax-snapshots-local + DataStax Local Snapshots + https://repo.aws.dsinternal.org/artifactory/datastax-snapshots-local/ + + diff --git a/src/main/java/org/apache/pulsar/io/jcloud/BlobStoreAbstractConfig.java b/src/main/java/org/apache/pulsar/io/jcloud/BlobStoreAbstractConfig.java index e1ce1518..fcda4f42 100644 --- a/src/main/java/org/apache/pulsar/io/jcloud/BlobStoreAbstractConfig.java +++ b/src/main/java/org/apache/pulsar/io/jcloud/BlobStoreAbstractConfig.java @@ -119,6 +119,7 @@ public class BlobStoreAbstractConfig implements Serializable { private String bytesFormatTypeSeparator = "0x10"; private boolean skipFailedMessages = false; private boolean jsonAllowNaN = false; + private String topicsToPathMapping; public void validate() { checkNotNull(provider, "provider not set."); @@ -194,6 +195,10 @@ public void validate() { if (parquetCodec != null && (parquetCodec.isEmpty() || parquetCodec.equals("none"))) { parquetCodec = null; } + + if (StringUtils.isNotEmpty(topicsToPathMapping)) { + validateTopicsToPathMapping(); + } } private static boolean hasURIScheme(String endpoint) { @@ -205,4 +210,20 @@ private static boolean hasURIScheme(String endpoint) { } } + public void validateTopicsToPathMapping() { + for (String topicToPathMapping : topicsToPathMapping.split(",")) { + String[] topicPathPair = topicToPathMapping.split("="); + checkArgument(topicPathPair.length == 2, + "Invalid topicsToPathMapping format: " + topicToPathMapping); + String topic = topicPathPair[0].trim(); + String path = topicPathPair[1].trim(); + checkArgument(isNotBlank(topic), "Topic in topicsToPathMapping cannot be empty."); + checkArgument(isNotBlank(path), "Path in topicsToPathMapping cannot be empty."); + checkArgument(topic.matches("^[^/]+/[^/]+/[^/]+$"), + "Topic must be in the format 'tenant/ns/topic': " + topic); + checkArgument(!path.startsWith("/") && !path.endsWith("/"), + "Path in topicsToPathMapping cannot starts or ends with '/': " + path); + } + } + } diff --git a/src/main/java/org/apache/pulsar/io/jcloud/format/BytesFormat.java b/src/main/java/org/apache/pulsar/io/jcloud/format/BytesFormat.java index 8b85b3c8..b6b1e39b 100644 --- a/src/main/java/org/apache/pulsar/io/jcloud/format/BytesFormat.java +++ b/src/main/java/org/apache/pulsar/io/jcloud/format/BytesFormat.java @@ -19,6 +19,7 @@ package org.apache.pulsar.io.jcloud.format; import java.nio.ByteBuffer; +import java.nio.charset.StandardCharsets; import java.util.Iterator; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.Schema; @@ -35,6 +36,7 @@ */ public class BytesFormat implements Format, InitConfiguration { + public static final byte[] NULL_BYTES = "".getBytes(StandardCharsets.UTF_8); private byte[] lineSeparatorBytes; @Override @@ -62,7 +64,10 @@ public ByteBuffer recordWriterBuf(Iterator> record) throws while (record.hasNext()) { final Record next = record.next(); final Message message = next.getMessage().get(); - final byte[] data = message.getData(); + byte[] data = message.getData(); + if (data == null) { + data = NULL_BYTES; + } dataOutput.write(data); dataOutput.write(lineSeparatorBytes); } diff --git a/src/main/java/org/apache/pulsar/io/jcloud/format/JsonFormat.java b/src/main/java/org/apache/pulsar/io/jcloud/format/JsonFormat.java index e310028d..f56a7e50 100644 --- a/src/main/java/org/apache/pulsar/io/jcloud/format/JsonFormat.java +++ b/src/main/java/org/apache/pulsar/io/jcloud/format/JsonFormat.java @@ -24,6 +24,7 @@ import com.fasterxml.jackson.core.JsonParser; import com.fasterxml.jackson.core.type.TypeReference; import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.exc.MismatchedInputException; import com.google.protobuf.DynamicMessage; @@ -128,11 +129,13 @@ public ByteBuffer recordWriterBuf(Iterator> record) throws return ByteBuffer.wrap(stringBuilder.toString().getBytes(StandardCharsets.UTF_8)); } - private Map convertRecordToObject(GenericRecord record, Schema schema) throws IOException { + protected Map convertRecordToObject(GenericRecord record, Schema schema) throws IOException { if (record.getSchemaType().isStruct()) { switch (record.getSchemaType()) { - case AVRO: case JSON: + JsonNode nativeObject = (JsonNode) record.getNativeObject(); + return JSON_MAPPER.get().convertValue(nativeObject, TYPEREF); + case AVRO: case PROTOBUF: { List fields = record.getFields(); Map result = new LinkedHashMap<>(fields.size()); @@ -142,6 +145,9 @@ private Map convertRecordToObject(GenericRecord record, Schema implements Partitioner { private boolean withTopicPartitionNumber; private boolean useIndexAsOffset; + private Map topicsToPathMapping; + private Map computedTopicsToPathMapping; + @Override public void configure(BlobStoreAbstractConfig config) { this.sliceTopicPartitionPath = config.isSliceTopicPartitionPath(); this.withTopicPartitionNumber = config.isWithTopicPartitionNumber(); this.useIndexAsOffset = config.isPartitionerUseIndexAsOffset(); + this.topicsToPathMapping = parseTopicsToPathMappingString(config.getTopicsToPathMapping()); + this.computedTopicsToPathMapping = new HashMap<>(); + } + + private Map parseTopicsToPathMappingString(String topicsToPathMappingString) { + Map topicsToPathMapping = new HashMap<>(); + if (StringUtils.isNotEmpty(topicsToPathMappingString)) { + for (String topicPathPair : topicsToPathMappingString.split(",")) { + String[] topicToPathMapping = topicPathPair.split("="); + if (topicToPathMapping.length == 2) { + topicsToPathMapping.put(topicToPathMapping[0], topicToPathMapping[1]); + } else { + LOGGER.warn("The topic to path mapping is not in correct format {}.", topicPathPair); + } + } + } + LOGGER.info("The final topics to path mapping is: {}", topicsToPathMapping); + return topicsToPathMapping; } @Override public String generatePartitionedPath(String topic, String encodedPartition) { + if (computedTopicsToPathMapping.containsKey(topic)) { + return String.format("%s%s%s", computedTopicsToPathMapping.get(topic), PATH_SEPARATOR, encodedPartition); + } else { + List joinList = new ArrayList<>(); + TopicName topicName = TopicName.get(topic); + joinList.add(topicName.getTenant()); + joinList.add(topicName.getNamespacePortion()); - List joinList = new ArrayList<>(); - TopicName topicName = TopicName.get(topic); - joinList.add(topicName.getTenant()); - joinList.add(topicName.getNamespacePortion()); - - if (topicName.isPartitioned() && withTopicPartitionNumber) { - if (sliceTopicPartitionPath) { + if (topicName.isPartitioned() && withTopicPartitionNumber) { + if (sliceTopicPartitionPath) { + TopicName newTopicName = TopicName.get(topicName.getPartitionedTopicName()); + joinList.add(newTopicName.getLocalName()); + joinList.add(Integer.toString(topicName.getPartitionIndex())); + } else { + joinList.add(topicName.getLocalName()); + } + } else { TopicName newTopicName = TopicName.get(topicName.getPartitionedTopicName()); joinList.add(newTopicName.getLocalName()); - joinList.add(Integer.toString(topicName.getPartitionIndex())); - } else { - joinList.add(topicName.getLocalName()); } - } else { - TopicName newTopicName = TopicName.get(topicName.getPartitionedTopicName()); - joinList.add(newTopicName.getLocalName()); + String generatedTopicPath = StringUtils.join(joinList, PATH_SEPARATOR); + if (topicsToPathMapping.containsKey(generatedTopicPath)) { + joinList.clear(); + joinList.add(topicsToPathMapping.get(generatedTopicPath)); + } + computedTopicsToPathMapping.put(topic, StringUtils.join(joinList, PATH_SEPARATOR)); + joinList.add(encodedPartition); + return StringUtils.join(joinList, PATH_SEPARATOR); } - joinList.add(encodedPartition); - return StringUtils.join(joinList, PATH_SEPARATOR); } protected long getMessageOffset(Record record) { diff --git a/src/main/java/org/apache/pulsar/io/jcloud/util/AvroRecordUtil.java b/src/main/java/org/apache/pulsar/io/jcloud/util/AvroRecordUtil.java index fea44b88..47949271 100644 --- a/src/main/java/org/apache/pulsar/io/jcloud/util/AvroRecordUtil.java +++ b/src/main/java/org/apache/pulsar/io/jcloud/util/AvroRecordUtil.java @@ -43,6 +43,7 @@ import org.apache.pulsar.client.impl.schema.generic.GenericAvroSchema; import org.apache.pulsar.client.impl.schema.generic.GenericJsonSchema; import org.apache.pulsar.client.impl.schema.generic.GenericProtobufNativeSchema; +import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.functions.api.Record; @@ -168,8 +169,10 @@ public static Schema convertToAvroSchema(org.apache.pulsar.client.api.Schema final Schema avroSchema = Schema.createRecord("KVSchema", null, null, false, Arrays.asList( - new Schema.Field("key", parseAvroSchema(keySchemaDef)), - new Schema.Field("value", parseAvroSchema(valueSchemaDef)) + // namespace+name must be different in the KVSchema + new Schema.Field("key", parseAvroSchema(keySchemaDef, "_key")), + new Schema.Field("value", Schema.createUnion( + Schema.create(Schema.Type.NULL), parseAvroSchema(valueSchemaDef, "_value"))) )); return avroSchema; } else { @@ -181,14 +184,24 @@ public static Schema convertToAvroSchema(org.apache.pulsar.client.api.Schema if (StringUtils.isEmpty(rootAvroSchemaString)) { throw new IllegalArgumentException("schema definition is empty"); } - return parseAvroSchema(rootAvroSchemaString); + return parseAvroSchema(rootAvroSchemaString, null); } } - private static Schema parseAvroSchema(String jsonSchema) { + private static Schema parseAvroSchema(String jsonSchema, String namespaceSuffix) { final Schema.Parser parser = new Schema.Parser(); parser.setValidateDefaults(false); - return parser.parse(jsonSchema); + Schema schema = parser.parse(jsonSchema); + List fields = schema.getFields() + .stream() + .map(f -> new Schema.Field(f, f.schema())) + .collect(Collectors.toList()); + if (namespaceSuffix == null) { + namespaceSuffix = ""; + } + final String namespace = schema.getNamespace() == null + ? namespaceSuffix : schema.getNamespace() + namespaceSuffix; + return Schema.createRecord(schema.getName(), null, namespace, false, fields); } public static org.apache.avro.generic.GenericRecord convertGenericRecord(DynamicMessage recordValue, @@ -215,6 +228,30 @@ public static org.apache.avro.generic.GenericRecord convertGenericRecord(Dynamic public static org.apache.avro.generic.GenericRecord convertGenericRecord(GenericRecord recordValue, Schema rootAvroSchema) { + if (recordValue.getSchemaType() == SchemaType.KEY_VALUE) { + KeyValue keyValue = + (KeyValue) recordValue.getNativeObject(); + org.apache.avro.generic.GenericRecord recordHolder = new GenericData.Record(rootAvroSchema); + GenericRecord keyObject = keyValue.getKey(); + if (keyObject != null) { + Schema keySchema = rootAvroSchema.getField("key").schema(); + recordHolder.put("key", convertGenericRecord(keyObject, keySchema)); + } + GenericRecord valueObject = keyValue.getValue(); + if (valueObject != null) { + Schema valueSchema = rootAvroSchema.getField("value").schema(); + recordHolder.put("value", convertGenericRecord(valueObject, valueSchema)); + } + return recordHolder; + } + + // handle nullable fields that are union[null,record] + if (rootAvroSchema.isUnion()) { + rootAvroSchema = rootAvroSchema.getTypes().stream() + .filter(schema -> schema.getType().equals(Schema.Type.RECORD)) + .findFirst() + .get(); + } org.apache.avro.generic.GenericRecord recordHolder = new GenericData.Record(rootAvroSchema); for (org.apache.pulsar.client.api.schema.Field field : recordValue.getFields()) { Schema.Field field1 = rootAvroSchema.getField(field.getName()); diff --git a/src/test/java/org/apache/pulsar/io/jcloud/ConnectorConfigTest.java b/src/test/java/org/apache/pulsar/io/jcloud/ConnectorConfigTest.java index ac55f914..e7e6ca8b 100644 --- a/src/test/java/org/apache/pulsar/io/jcloud/ConnectorConfigTest.java +++ b/src/test/java/org/apache/pulsar/io/jcloud/ConnectorConfigTest.java @@ -52,6 +52,7 @@ public void loadBasicConfigTest() throws IOException { config.put("timePartitionDuration", "2d"); config.put("batchSize", 10); config.put("partitioner", "topic"); + config.put("topicsToPathMapping", "my-tenant/my-ns/topic1=path-mytopic1,my-tenant/my-ns/topic2=path-mytopic2"); CloudStorageSinkConfig cloudStorageSinkConfig = CloudStorageSinkConfig.load(config); cloudStorageSinkConfig.validate(); @@ -69,6 +70,7 @@ public void loadBasicConfigTest() throws IOException { cloudStorageSinkConfig.getPartitioner().toString().toLowerCase()); Assert.assertEquals((int) config.get("batchSize"), cloudStorageSinkConfig.getPendingQueueSize()); Assert.assertEquals(10000000L, cloudStorageSinkConfig.getMaxBatchBytes()); + Assert.assertEquals(config.get("topicsToPathMapping"), cloudStorageSinkConfig.getTopicsToPathMapping()); } @Test @@ -131,6 +133,77 @@ public void timePartitionDurationTest() throws IOException { } } + @Test + public void topicsToPathMappingTest() throws IOException { + Map config = new HashMap<>(); + config.put("provider", PROVIDER_AWSS3); + config.put("accessKeyId", "aws-s3"); + config.put("secretAccessKey", "aws-s3"); + config.put("bucket", "testbucket"); + config.put("region", "localhost"); + config.put("endpoint", "https://us-standard"); + config.put("formatType", "avro"); + config.put("partitionerType", "default"); + config.put("timePartitionPattern", "yyyy-MM-dd"); + config.put("timePartitionDuration", "2d"); + config.put("batchSize", 10); + try { + config.put("topicsToPathMapping", "tenant1/ns1/topic1=path-mytopic1,tenant1/ns1/topic2=path-mytopic2"); + CloudStorageSinkConfig.load(config).validate(); + } catch (Exception e) { + Assert.fail(); + } + + try { + config.put("topicsToPathMapping", "topic1=path-mytopic1"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/=path-mytopic1"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/ns1=path-mytopic1"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/ns1/topic1=path-mytopic1,tenant1/ns1/topic2="); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/ns1/topic1=path-mytopic1,=path-mytopic2"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/ns1/topic1=/path-mytopic1"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + + try { + config.put("topicsToPathMapping", "tenant1/ns1/topic1=path-mytopic1/"); + CloudStorageSinkConfig.load(config).validate(); + Assert.fail(); + } catch (Exception e) { + } + } + @Test public void pathPrefixTest() throws IOException { Map config = new HashMap<>(); diff --git a/src/test/java/org/apache/pulsar/io/jcloud/PulsarTestBase.java b/src/test/java/org/apache/pulsar/io/jcloud/PulsarTestBase.java index 1c8ec760..da544261 100644 --- a/src/test/java/org/apache/pulsar/io/jcloud/PulsarTestBase.java +++ b/src/test/java/org/apache/pulsar/io/jcloud/PulsarTestBase.java @@ -37,6 +37,7 @@ import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SubscriptionInitialPosition; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.schema.KeyValue; import org.apache.pulsar.common.schema.SchemaType; @@ -75,7 +76,7 @@ public static void prepare() throws Exception { log.info("-------------------------------------------------------------------------"); - final String pulsarImage = System.getProperty("pulsar.systemtest.image", "streamnative/pulsar:2.9.2.13"); + final String pulsarImage = System.getProperty("pulsar.systemtest.image", "apachepulsar/pulsar:latest"); pulsarService = new PulsarContainer(DockerImageName.parse(pulsarImage)); pulsarService.waitingFor(new HttpWaitStrategy() .forPort(BROKER_HTTP_PORT) @@ -117,7 +118,6 @@ public static List sendTypedMessages( return sendTypedMessages(topic, type, messages, partition, null); } - @SuppressWarnings("unchecked") public static List sendTypedMessages( String topic, SchemaType type, @@ -125,6 +125,66 @@ public static List sendTypedMessages( Optional partition, Class tClass) throws PulsarClientException { + Schema schema; + + switch (type) { + case BOOLEAN: + schema = Schema.BOOL; + break; + case BYTES: + schema = Schema.BYTES; + break; + case DATE: + schema = Schema.DATE; + break; + case STRING: + schema = Schema.STRING; + break; + case TIMESTAMP: + schema = Schema.TIMESTAMP; + break; + case INT8: + schema = Schema.INT8; + break; + case DOUBLE: + schema = Schema.DOUBLE; + break; + case FLOAT: + schema = Schema.FLOAT; + break; + case INT32: + schema = Schema.INT32; + break; + case INT16: + schema = Schema.INT16; + break; + case INT64: + schema = Schema.INT64; + break; + case AVRO: + schema = Schema.AVRO(tClass); + break; + case JSON: + schema = Schema.JSON(tClass); + break; + case KEY_VALUE: + final KeyValue kvMessage = (KeyValue) messages.get(0); + schema = Schema.KeyValue(kvMessage.getKey().getClass(), + kvMessage.getValue().getClass()); + break; + default: + throw new NotImplementedException("Unsupported type " + type); + } + return sendTypedMessages(topic, messages, schema, partition); + + } + @SuppressWarnings("unchecked") + public static List sendTypedMessages( + String topic, + List messages, + Schema produceSchema, + Optional partition) throws PulsarClientException { + String topicName; if (partition.isPresent()) { topicName = topic + TopicName.PARTITIONED_TOPIC_SUFFIX + partition.get(); @@ -138,56 +198,7 @@ public static List sendTypedMessages( try { client = PulsarClient.builder().serviceUrl(getServiceUrl()).build(); - - switch (type) { - case BOOLEAN: - producer = (Producer) client.newProducer(Schema.BOOL).topic(topicName).create(); - break; - case BYTES: - producer = (Producer) client.newProducer(Schema.BYTES).topic(topicName).create(); - break; - case DATE: - producer = (Producer) client.newProducer(Schema.DATE).topic(topicName).create(); - break; - case STRING: - producer = (Producer) client.newProducer(Schema.STRING).topic(topicName).create(); - break; - case TIMESTAMP: - producer = (Producer) client.newProducer(Schema.TIMESTAMP).topic(topicName).create(); - break; - case INT8: - producer = (Producer) client.newProducer(Schema.INT8).topic(topicName).create(); - break; - case DOUBLE: - producer = (Producer) client.newProducer(Schema.DOUBLE).topic(topicName).create(); - break; - case FLOAT: - producer = (Producer) client.newProducer(Schema.FLOAT).topic(topicName).create(); - break; - case INT32: - producer = (Producer) client.newProducer(Schema.INT32).topic(topicName).create(); - break; - case INT16: - producer = (Producer) client.newProducer(Schema.INT16).topic(topicName).create(); - break; - case INT64: - producer = (Producer) client.newProducer(Schema.INT64).topic(topicName).create(); - break; - case AVRO: - producer = (Producer) client.newProducer(Schema.AVRO(tClass)).topic(topicName).create(); - break; - case JSON: - producer = (Producer) client.newProducer(Schema.JSON(tClass)).topic(topicName).create(); - break; - case KEY_VALUE: - final KeyValue kvMessage = (KeyValue) messages.get(0); - producer = (Producer) client.newProducer(Schema.KeyValue(kvMessage.getKey().getClass(), - kvMessage.getValue().getClass())).topic(topicName).create(); - break; - - default: - throw new NotImplementedException("Unsupported type " + type); - } + producer = client.newProducer(produceSchema).topic(topicName).create(); for (T message : messages) { MessageId mid = producer.send(message); @@ -264,6 +275,7 @@ public static void consumerMessages(String topic, consumer = client.newConsumer(schema) .topic(topic) + .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) .subscriptionName("test") .subscribe(); int receiveCount = 0; diff --git a/src/test/java/org/apache/pulsar/io/jcloud/bo/TestRecord.java b/src/test/java/org/apache/pulsar/io/jcloud/bo/TestRecord.java index 158170a5..726b5875 100644 --- a/src/test/java/org/apache/pulsar/io/jcloud/bo/TestRecord.java +++ b/src/test/java/org/apache/pulsar/io/jcloud/bo/TestRecord.java @@ -40,6 +40,7 @@ public class TestRecord { @AllArgsConstructor @NoArgsConstructor public static class TestSubRecord { + public static final String FILE_NAME = "name"; private String name; } } diff --git a/src/test/java/org/apache/pulsar/io/jcloud/container/PulsarContainer.java b/src/test/java/org/apache/pulsar/io/jcloud/container/PulsarContainer.java index 31804099..4c67b6fb 100644 --- a/src/test/java/org/apache/pulsar/io/jcloud/container/PulsarContainer.java +++ b/src/test/java/org/apache/pulsar/io/jcloud/container/PulsarContainer.java @@ -32,9 +32,9 @@ public class PulsarContainer extends GenericContainer { public static final int BROKER_HTTP_PORT = 8080; public static final String METRICS_ENDPOINT = "/metrics"; - private static final DockerImageName DEFAULT_IMAGE_NAME = DockerImageName.parse("streamnative/pulsar"); + private static final DockerImageName DEFAULT_IMAGE_NAME = DockerImageName.parse("apachepulsar/pulsar"); @Deprecated - private static final String DEFAULT_TAG = "2.8.1.6"; + private static final String DEFAULT_TAG = "latest"; private boolean functionsWorkerEnabled = false; @@ -57,7 +57,7 @@ public PulsarContainer(String pulsarVersion) { public PulsarContainer(final DockerImageName dockerImageName) { super(dockerImageName); - dockerImageName.assertCompatibleWith(DockerImageName.parse("streamnative/pulsar")); + dockerImageName.assertCompatibleWith(DockerImageName.parse("apachepulsar/pulsar")); withExposedPorts(BROKER_PORT, BROKER_HTTP_PORT); withCommand("/pulsar/bin/pulsar", "standalone", "--no-functions-worker", "-nss"); diff --git a/src/test/java/org/apache/pulsar/io/jcloud/format/FormatTestBase.java b/src/test/java/org/apache/pulsar/io/jcloud/format/FormatTestBase.java index 035703bd..d9621d13 100644 --- a/src/test/java/org/apache/pulsar/io/jcloud/format/FormatTestBase.java +++ b/src/test/java/org/apache/pulsar/io/jcloud/format/FormatTestBase.java @@ -29,6 +29,8 @@ import java.util.Optional; import java.util.function.Consumer; import java.util.stream.Collectors; +import org.apache.avro.SchemaBuilder; +import org.apache.avro.generic.GenericData; import org.apache.avro.util.Utf8; import org.apache.commons.lang3.RandomStringUtils; import org.apache.pulsar.client.admin.PulsarAdmin; @@ -37,10 +39,12 @@ import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.schema.Field; import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.SchemaDefinition; import org.apache.pulsar.client.impl.schema.AutoConsumeSchema; import org.apache.pulsar.client.impl.schema.generic.GenericJsonRecord; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.schema.KeyValue; +import org.apache.pulsar.common.schema.KeyValueEncodingType; import org.apache.pulsar.common.schema.SchemaType; import org.apache.pulsar.io.jcloud.BlobStoreAbstractConfig; import org.apache.pulsar.io.jcloud.PulsarTestBase; @@ -72,6 +76,8 @@ public abstract class FormatTestBase extends PulsarTestBase { TopicName.get("test-json-bytes-parquet-json" + RandomStringUtils.randomAlphabetic(5)); protected static TopicName jsonStringTopicName = TopicName.get("test-json-string-parquet-json" + RandomStringUtils.randomAlphabetic(5)); + private static final TopicName kvSeparatedTopicName = + TopicName.get("test-parquet-kv-sep" + RandomStringUtils.randomAlphabetic(5)); @BeforeClass public static void setUp() throws Exception { @@ -169,6 +175,61 @@ public void testProtobufNativeRecordWriter() throws Exception { consumerMessages(protobufNativeTopicName.toString(), Schema.AUTO_CONSUME(), handle, testRecords.size(), 2000); } + @Test + public void testKeyValueAvroWithSameNamespaceName() throws Exception { + org.apache.avro.Schema keySchema = SchemaBuilder.record("record") + .fields() + .name("id").type().stringType().noDefault() + .endRecord(); + org.apache.avro.Schema valueSchema = SchemaBuilder.record("record") + .fields() + .name("content").type().stringType().noDefault() + .endRecord(); + + SchemaDefinition keySchemaDef = SchemaDefinition.builder() + .withJsonDef(String.format("" + + "{\n" + + " \t\"type\": \"record\",\n" + + " \t\"name\": \"record\",\n" + + " \t\"fields\": [{\n" + + " \t\t\"name\": \"id\",\n" + + " \t\t\"type\": [\"string\"]\n" + + " \t}]\n" + + " }")) + .build(); + + SchemaDefinition valueSchemaDef = SchemaDefinition.builder() + .withJsonDef(String.format("" + + "{\n" + + " \t\"type\": \"record\",\n" + + " \t\"name\": \"record\",\n" + + " \t\"fields\": [{\n" + + " \t\t\"name\": \"content\",\n" + + " \t\t\"type\": [\"string\"]\n" + + " \t}]\n" + + " }")) + .build(); + + + Schema> schema = Schema.KeyValue(Schema.AVRO(keySchemaDef), + Schema.AVRO(valueSchemaDef), KeyValueEncodingType.SEPARATED); + + GenericData.Record keyRecord = new GenericData.Record(keySchema); + keyRecord.put("id", "theid"); + GenericData.Record valueRecord = new GenericData.Record(valueSchema); + valueRecord.put("content", "ccc"); + List testRecords = Arrays.asList( + new KeyValue(keyRecord, valueRecord), + new KeyValue(keyRecord, null) + ); + + sendTypedMessages(kvSeparatedTopicName.toString(), testRecords, schema, Optional.empty()); + + Consumer> handle = getMessageConsumer(kvSeparatedTopicName); + consumerMessages(kvSeparatedTopicName.toString(), Schema.AUTO_CONSUME(), handle, testRecords.size(), 2000); + } + + protected abstract boolean supportMetadata(); protected Consumer> getMessageConsumer(TopicName topic) { @@ -224,11 +285,11 @@ public abstract org.apache.avro.generic.GenericRecord getFormatGeneratedRecord(T throws Exception; public abstract DynamicMessage getDynamicMessage(TopicName topicName, - Message msg) + Message msg) throws Exception; public abstract Map getJSONMessage(TopicName topicName, - Message msg) + Message msg) throws Exception; protected void assertEquals(DynamicMessage msgValue, org.apache.avro.generic.GenericRecord record) { @@ -237,16 +298,16 @@ protected void assertEquals(DynamicMessage msgValue, org.apache.avro.generic.Gen Object sourceValue = msgValue.getField(descriptor.findFieldByName(field.name())); Object newValue = record.get(field.name()); Assert.assertEquals( - MessageFormat.format( - "field[{0} sourceValue [{1}:{2}] not equal newValue [{3}:{4}]", - field.name(), - sourceValue, - sourceValue != null ? sourceValue.getClass().getName() : "null", - newValue, - newValue != null ? newValue.getClass().getName() : "null" - ), - sourceValue, - newValue); + MessageFormat.format( + "field[{0} sourceValue [{1}:{2}] not equal newValue [{3}:{4}]", + field.name(), + sourceValue, + sourceValue != null ? sourceValue.getClass().getName() : "null", + newValue, + newValue != null ? newValue.getClass().getName() : "null" + ), + sourceValue, + newValue); } } @@ -382,16 +443,16 @@ protected void assertEquals(DynamicMessage msgValue, Map record) Object newValue = record.get(fieldName); if (!fieldDescriptor.isRepeated()) { Assert.assertEquals( - MessageFormat.format( - "field[{0} sourceValue [{1}:{2}] not equal newValue [{3}:{4}]", - fieldName, - sourceValue, - sourceValue != null ? sourceValue.getClass().getName() : "null", - newValue, - newValue != null ? newValue.getClass().getName() : "null" - ), - sourceValue, - newValue); + MessageFormat.format( + "field[{0} sourceValue [{1}:{2}] not equal newValue [{3}:{4}]", + fieldName, + sourceValue, + sourceValue != null ? sourceValue.getClass().getName() : "null", + newValue, + newValue != null ? newValue.getClass().getName() : "null" + ), + sourceValue, + newValue); } } } diff --git a/src/test/java/org/apache/pulsar/io/jcloud/format/JsonFormatMethodTest.java b/src/test/java/org/apache/pulsar/io/jcloud/format/JsonFormatMethodTest.java new file mode 100644 index 00000000..208a5367 --- /dev/null +++ b/src/test/java/org/apache/pulsar/io/jcloud/format/JsonFormatMethodTest.java @@ -0,0 +1,66 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.io.jcloud.format; + +import com.fasterxml.jackson.annotation.JsonInclude; +import com.fasterxml.jackson.databind.DeserializationFeature; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.util.HashMap; +import java.util.Map; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.GenericRecordBuilder; +import org.apache.pulsar.client.impl.schema.generic.GenericJsonSchema; +import org.apache.pulsar.io.jcloud.BlobStoreAbstractConfig; +import org.apache.pulsar.io.jcloud.bo.TestRecord; +import org.junit.Assert; +import org.junit.Test; + +public class JsonFormatMethodTest { + + private static final ThreadLocal JSON_MAPPER = ThreadLocal.withInitial(() -> { + ObjectMapper mapper = new ObjectMapper(); + mapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); + mapper.setSerializationInclusion(JsonInclude.Include.NON_NULL); + return mapper; + }); + + @Test + public void testJsonIgnoreSchemaRead() throws Exception { + + // 1. Gen GenericRecord with TestSubRecord schema but set to incompatible data. + Schema subRecordSchema = Schema.JSON(TestRecord.TestSubRecord.class); + GenericJsonSchema genericRecordSchema = new GenericJsonSchema(subRecordSchema.getSchemaInfo()); + GenericRecordBuilder genericRecordBuilder = genericRecordSchema.newRecordBuilder(); + Map originalData = new HashMap<>(); + originalData.put("incompatibleKey1", "value1"); + originalData.put("incompatibleKey2", 2); + originalData.forEach(genericRecordBuilder::set); + GenericRecord genericRecord = genericRecordBuilder.build(); + + // 3. Will be converted to the same JSON string as the original data. + JsonFormat jsonFormat = new JsonFormat(); + BlobStoreAbstractConfig blobStoreAbstractConfig = new BlobStoreAbstractConfig(); + jsonFormat.configure(blobStoreAbstractConfig); + Map jsonMap = jsonFormat.convertRecordToObject(genericRecord, genericRecordSchema); + String convertedJsonString = JSON_MAPPER.get().writeValueAsString(jsonMap); + String originalJsonString = JSON_MAPPER.get().writeValueAsString(originalData); + Assert.assertEquals(originalJsonString, convertedJsonString); + } +} diff --git a/src/test/java/org/apache/pulsar/io/jcloud/partitioner/TopicsToPathMappingPartitionerTest.java b/src/test/java/org/apache/pulsar/io/jcloud/partitioner/TopicsToPathMappingPartitionerTest.java new file mode 100644 index 00000000..2044a655 --- /dev/null +++ b/src/test/java/org/apache/pulsar/io/jcloud/partitioner/TopicsToPathMappingPartitionerTest.java @@ -0,0 +1,260 @@ +/** + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ +package org.apache.pulsar.io.jcloud.partitioner; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; +import com.google.common.base.Supplier; +import java.io.File; +import java.text.MessageFormat; +import java.util.Optional; +import junit.framework.TestCase; +import org.apache.pulsar.client.api.Message; +import org.apache.pulsar.client.impl.MessageIdImpl; +import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.functions.api.Record; +import org.apache.pulsar.io.jcloud.BlobStoreAbstractConfig; +import org.apache.pulsar.io.jcloud.partitioner.legacy.Partitioner; +import org.apache.pulsar.io.jcloud.partitioner.legacy.SimplePartitioner; +import org.apache.pulsar.io.jcloud.partitioner.legacy.TimePartitioner; +import org.junit.Assert; +import org.junit.Test; +import org.junit.runner.RunWith; +import org.junit.runners.Parameterized; + + +/** + * partitioner unit test. + */ +@RunWith(Parameterized.class) +public class TopicsToPathMappingPartitionerTest extends TestCase { + + private static final String pathSeparator = File.separator; + + @Parameterized.Parameter(0) + public Partitioner partitioner; + + @Parameterized.Parameter(1) + public String expected; + + @Parameterized.Parameter(2) + public String expectedPartitionedPath; + + @Parameterized.Parameter(3) + public Record pulsarRecord; + + @Parameterized.Parameters + public static Object[][] data() { + BlobStoreAbstractConfig blobStoreAbstractConfig1 = new BlobStoreAbstractConfig(); + blobStoreAbstractConfig1.setTimePartitionDuration("1d"); + blobStoreAbstractConfig1.setTimePartitionPattern("yyyy-MM-dd"); + blobStoreAbstractConfig1.setSliceTopicPartitionPath(false); + blobStoreAbstractConfig1.setWithTopicPartitionNumber(false); + blobStoreAbstractConfig1.setTopicsToPathMapping("public/default/test=path1," + + "public/default/test-partition-1=path2,public/default/test/1=path3"); + SimplePartitioner simplePartitioner1 = new SimplePartitioner<>(); + simplePartitioner1.configure(blobStoreAbstractConfig1); + TimePartitioner dayPartitioner1 = new TimePartitioner<>(); + dayPartitioner1.configure(blobStoreAbstractConfig1); + + BlobStoreAbstractConfig blobStoreAbstractConfig2 = new BlobStoreAbstractConfig(); + blobStoreAbstractConfig2.setTimePartitionDuration("1d"); + blobStoreAbstractConfig2.setTimePartitionPattern("yyyy-MM-dd"); + blobStoreAbstractConfig2.setSliceTopicPartitionPath(false); + blobStoreAbstractConfig2.setWithTopicPartitionNumber(true); + blobStoreAbstractConfig2.setTopicsToPathMapping("public/default/test=path1," + + "public/default/test-partition-1=path2,public/default/test/1=path3"); + SimplePartitioner simplePartitioner2 = new SimplePartitioner<>(); + simplePartitioner2.configure(blobStoreAbstractConfig2); + TimePartitioner dayPartitioner2 = new TimePartitioner<>(); + dayPartitioner2.configure(blobStoreAbstractConfig2); + + BlobStoreAbstractConfig blobStoreAbstractConfig3 = new BlobStoreAbstractConfig(); + blobStoreAbstractConfig3.setTimePartitionDuration("1d"); + blobStoreAbstractConfig3.setTimePartitionPattern("yyyy-MM-dd"); + blobStoreAbstractConfig3.setSliceTopicPartitionPath(true); + blobStoreAbstractConfig3.setWithTopicPartitionNumber(false); + blobStoreAbstractConfig3.setTopicsToPathMapping("public/default/test=path1," + + "public/default/test-partition-1=path2,public/default/test/1=path3"); + SimplePartitioner simplePartitioner3 = new SimplePartitioner<>(); + simplePartitioner3.configure(blobStoreAbstractConfig3); + TimePartitioner dayPartitioner3 = new TimePartitioner<>(); + dayPartitioner3.configure(blobStoreAbstractConfig3); + + BlobStoreAbstractConfig blobStoreAbstractConfig4 = new BlobStoreAbstractConfig(); + blobStoreAbstractConfig4.setTimePartitionDuration("1d"); + blobStoreAbstractConfig4.setTimePartitionPattern("yyyy-MM-dd"); + blobStoreAbstractConfig4.setSliceTopicPartitionPath(true); + blobStoreAbstractConfig4.setWithTopicPartitionNumber(true); + blobStoreAbstractConfig4.setTopicsToPathMapping("public/default/test=path1," + + "public/default/test-partition-1=path2,public/default/test/1=path3"); + SimplePartitioner simplePartitioner4 = new SimplePartitioner<>(); + simplePartitioner4.configure(blobStoreAbstractConfig4); + TimePartitioner dayPartitioner4 = new TimePartitioner<>(); + dayPartitioner4.configure(blobStoreAbstractConfig4); + + return new Object[][]{ + new Object[]{ + simplePartitioner1, + "3221225506", + "path1" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + dayPartitioner1, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + simplePartitioner1, + "3221225506", + "path1" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + dayPartitioner1, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + simplePartitioner2, + "3221225506", + "path1" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + dayPartitioner2, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + simplePartitioner2, + "3221225506", + "path2" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + dayPartitioner2, + "2020-09-08" + pathSeparator + "3221225506", + "path2/2020-09-08" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + simplePartitioner3, + "3221225506", + "path1" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + dayPartitioner3, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + simplePartitioner3, + "3221225506", + "path1" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + dayPartitioner3, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + simplePartitioner4, + "3221225506", + "path1" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + dayPartitioner4, + "2020-09-08" + pathSeparator + "3221225506", + "path1/2020-09-08" + pathSeparator + "3221225506", + getTopic() + }, + new Object[]{ + simplePartitioner4, + "3221225506", + "path3" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + new Object[]{ + dayPartitioner4, + "2020-09-08" + pathSeparator + "3221225506", + "path3/2020-09-08" + pathSeparator + "3221225506", + getPartitionedTopic() + }, + }; + } + + public static Record getPartitionedTopic() { + @SuppressWarnings("unchecked") + Message mock = mock(Message.class); + when(mock.getPublishTime()).thenReturn(1599578218610L); + when(mock.getMessageId()).thenReturn(new MessageIdImpl(12, 34, 1)); + String topic = TopicName.get("test-partition-1").toString(); + Record mockRecord = mock(Record.class); + when(mockRecord.getTopicName()).thenReturn(Optional.of(topic)); + when(mockRecord.getPartitionIndex()).thenReturn(Optional.of(1)); + when(mockRecord.getMessage()).thenReturn(Optional.of(mock)); + when(mockRecord.getPartitionId()).thenReturn(Optional.of(String.format("%s-%s", topic, 1))); + when(mockRecord.getRecordSequence()).thenReturn(Optional.of(3221225506L)); + return mockRecord; + } + + public static Record getTopic() { + @SuppressWarnings("unchecked") + Message mock = mock(Message.class); + when(mock.getPublishTime()).thenReturn(1599578218610L); + when(mock.getMessageId()).thenReturn(new MessageIdImpl(12, 34, 1)); + String topic = TopicName.get("test").toString(); + Record mockRecord = mock(Record.class); + when(mockRecord.getTopicName()).thenReturn(Optional.of(topic)); + when(mockRecord.getPartitionIndex()).thenReturn(Optional.of(1)); + when(mockRecord.getMessage()).thenReturn(Optional.of(mock)); + when(mockRecord.getPartitionId()).thenReturn(Optional.of(String.format("%s-%s", topic, 1))); + when(mockRecord.getRecordSequence()).thenReturn(Optional.of(3221225506L)); + return mockRecord; + } + + @Test + public void testEncodePartition() { + String encodePartition = partitioner.encodePartition(pulsarRecord, System.currentTimeMillis()); + Supplier supplier = + () -> MessageFormat.format("expected: {0}\nactual: {1}", expected, encodePartition); + Assert.assertEquals(supplier.get(), expected, encodePartition); + } + + @Test + public void testGeneratePartitionedPath() { + String encodePartition = partitioner.encodePartition(pulsarRecord, System.currentTimeMillis()); + String partitionedPath = + partitioner.generatePartitionedPath(pulsarRecord.getTopicName().get(), encodePartition); + + Supplier supplier = + () -> MessageFormat.format("expected: {0}\nactual: {1}", expected, encodePartition); + Assert.assertEquals(supplier.get(), expectedPartitionedPath, partitionedPath); + } +}