diff --git a/Dockerfile b/Dockerfile
index f5d251e9..736a6012 100644
--- a/Dockerfile
+++ b/Dockerfile
@@ -1,59 +1,37 @@
-FROM maven:3.9.4-eclipse-temurin-11-focal AS build-core
+FROM public.ecr.aws/docker/library/maven:3.9.4-eclipse-temurin-11-focal AS build-core
COPY . /app
RUN mvn clean install -DskipTests -f /app/pom.xml
-# RUN mvn clean install -DskipTests -f /app/dataset-registry/pom.xml
-# RUN mvn clean install -DskipTests -f /app/transformation-sdk/pom.xml
-FROM maven:3.9.4-eclipse-temurin-11-focal AS build-pipeline
+FROM public.ecr.aws/docker/library/maven:3.9.4-eclipse-temurin-11-focal AS build-pipeline
COPY --from=build-core /root/.m2 /root/.m2
COPY . /app
RUN mvn clean package -DskipTests -f /app/pipeline/pom.xml
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS extractor-image
+FROM public.ecr.aws/docker/library/flink:1.20-scala_2.12-java11 AS unified-image
USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
-COPY --from=build-pipeline /app/pipeline/extractor/target/extractor-1.0.0.jar $FLINK_HOME/usrlib/
-
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS preprocessor-image
-USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
-COPY --from=build-pipeline /app/pipeline/preprocessor/target/preprocessor-1.0.0.jar $FLINK_HOME/usrlib/
-
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS denormalizer-image
-USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
-COPY --from=build-pipeline /app/pipeline/denormalizer/target/denormalizer-1.0.0.jar $FLINK_HOME/usrlib/
-
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS transformer-image
-USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
-COPY --from=build-pipeline /app/pipeline/transformer/target/transformer-1.0.0.jar $FLINK_HOME/usrlib/
-
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS dataset-router-image
-USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
-COPY --from=build-pipeline /app/pipeline/dataset-router/target/dataset-router-1.0.0.jar $FLINK_HOME/usrlib/
-
-# unified image build
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS unified-image
-USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
+# Move the bundled flink-s3-fs-hadoop plugin from opt/ to the required plugins subfolder.
+# This avoids a network download and guarantees the plugin version matches the runtime.
+RUN mkdir -p $FLINK_HOME/usrlib && \
+ mkdir -p $FLINK_HOME/plugins/flink-s3-fs-hadoop && \
+ mv $FLINK_HOME/opt/flink-s3-fs-hadoop-*.jar $FLINK_HOME/plugins/flink-s3-fs-hadoop/
+# Use IRSA/OIDC (Web Identity Token) for S3 auth instead of static access keys.
+# EKS injects AWS_ROLE_ARN and AWS_WEB_IDENTITY_TOKEN_FILE into pods whose service
+# account has an IAM role annotation; WebIdentityTokenCredentialsProvider reads them.
+RUN if [ -f "$FLINK_HOME/conf/config.yaml" ]; then \
+ echo 's3.aws.credentials.provider: com.amazonaws.auth.WebIdentityTokenCredentialsProvider' >> $FLINK_HOME/conf/config.yaml; \
+ else \
+ echo 's3.aws.credentials.provider: com.amazonaws.auth.WebIdentityTokenCredentialsProvider' >> $FLINK_HOME/conf/flink-conf.yaml; \
+ fi
COPY --from=build-pipeline /app/pipeline/unified-pipeline/target/unified-pipeline-1.0.0.jar $FLINK_HOME/usrlib/
-# # Lakehouse connector image build
-# FROM sanketikahub/flink:1.17.2-scala_2.12-java11 AS lakehouse-connector-image
-# USER flink
-# RUN wget https://repo1.maven.org/maven2/org/apache/flink/flink-shaded-hadoop-2-uber/2.8.3-10.0/flink-shaded-hadoop-2-uber-2.8.3-10.0.jar
-# RUN wget https://repo1.maven.org/maven2/org/apache/flink/flink-s3-fs-hadoop/1.17.2/flink-s3-fs-hadoop-1.17.2.jar
-# RUN wget https://repo.maven.apache.org/maven2/org/apache/hudi/hudi-flink1.17-bundle/1.0.2/hudi-flink1.17-bundle-1.0.2.jar
-# RUN mv flink-shaded-hadoop-2-uber-2.8.3-10.0.jar $FLINK_HOME/lib
-# RUN mv flink-s3-fs-hadoop-1.17.2.jar $FLINK_HOME/lib
-# RUN mv hudi-flink1.17-bundle-1.0.2.jar $FLINK_HOME/lib
-# # RUN mkdir $FLINK_HOME/custom-lib
-# COPY --from=build-pipeline /app/pipeline/hudi-connector/target/hudi-connector-1.0.0.jar $FLINK_HOME/lib
-
-# cache indexer image build
-FROM sanketikahub/flink:1.20-scala_2.12-java11 AS cache-indexer-image
+FROM public.ecr.aws/docker/library/flink:1.20-scala_2.12-java11 AS cache-indexer-image
USER flink
-RUN mkdir -p $FLINK_HOME/usrlib
+RUN mkdir -p $FLINK_HOME/usrlib && \
+ mkdir -p $FLINK_HOME/plugins/flink-s3-fs-hadoop && \
+ mv $FLINK_HOME/opt/flink-s3-fs-hadoop-*.jar $FLINK_HOME/plugins/flink-s3-fs-hadoop/
+RUN if [ -f "$FLINK_HOME/conf/config.yaml" ]; then \
+ echo 's3.aws.credentials.provider: com.amazonaws.auth.WebIdentityTokenCredentialsProvider' >> $FLINK_HOME/conf/config.yaml; \
+ else \
+ echo 's3.aws.credentials.provider: com.amazonaws.auth.WebIdentityTokenCredentialsProvider' >> $FLINK_HOME/conf/flink-conf.yaml; \
+ fi
COPY --from=build-pipeline /app/pipeline/cache-indexer/target/cache-indexer-1.0.0.jar $FLINK_HOME/usrlib/
\ No newline at end of file
diff --git a/FlinkDockerfile b/FlinkDockerfile
index b44ef9bd..8b5b0d33 100644
--- a/FlinkDockerfile
+++ b/FlinkDockerfile
@@ -1,4 +1,4 @@
-FROM --platform=linux/x86_64 flink:1.20-scala_2.12-java11
+FROM --platform=linux/x86_64 public.ecr.aws/docker/library/flink:1.20-scala_2.12-java11
USER flink
# RUN mkdir $FLINK_HOME/custom-lib
RUN wget https://repo1.maven.org/maven2/org/apache/flink/flink-shaded-hadoop-2-uber/2.8.3-10.0/flink-shaded-hadoop-2-uber-2.8.3-10.0.jar
diff --git a/dataset-registry/pom.xml b/dataset-registry/pom.xml
index a6944775..cf1e2284 100644
--- a/dataset-registry/pom.xml
+++ b/dataset-registry/pom.xml
@@ -34,7 +34,18 @@
org.apache.httpcomponents
httpclient
- 4.5.1
+ 4.5.14
+
+
+ commons-codec
+ commons-codec
+
+
+
+
+ commons-codec
+ commons-codec
+ 1.15
com.google.code.gson
@@ -55,7 +66,7 @@
junit
junit
- 4.12
+ 4.13.2
test
@@ -63,6 +74,32 @@
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.apache.commons
+ commons-compress
+
+
+ commons-io
+ commons-io
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
+ test
com.github.codemonstur
@@ -89,11 +126,45 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
+ test
org.apache.flink
flink-connector-kafka
3.3.0-1.20
+
+
+ org.apache.kafka
+ kafka-clients
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
@@ -164,7 +235,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
diff --git a/framework/pom.xml b/framework/pom.xml
index 6b973a7b..6d57d98f 100644
--- a/framework/pom.xml
+++ b/framework/pom.xml
@@ -39,6 +39,25 @@
org.apache.flink
flink-streaming-scala_${scala.maj.version}
${flink.version}
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.apache.flink
@@ -58,6 +77,14 @@
com.fasterxml.jackson.core
jackson-core
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
@@ -89,6 +116,28 @@
org.apache.kafka
kafka-clients
${kafka.version}
+
+
+ org.lz4
+ lz4-java
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ at.yawk.lz4
+ lz4-java
+ 1.10.3
+
+
+
+ org.xerial.snappy
+ snappy-java
+ 1.1.10.5
joda-time
@@ -99,7 +148,18 @@
org.apache.httpcomponents
httpclient
4.5.13
+
+
+ commons-codec
+ commons-codec
+
+
+
+ commons-codec
+ commons-codec
+ 1.15
+
com.google.code.gson
gson
@@ -142,7 +202,7 @@
org.postgresql
postgresql
- 42.6.0
+ 42.7.7
@@ -160,32 +220,199 @@
junit
junit
- 4.13.1
+ 4.13.2
test
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.apache.flink
flink-test-utils
${flink.version}
test
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.logging.log4j
+ log4j-api
+ ${log4j.version}
+ test
+
+
+ org.apache.logging.log4j
+ log4j-core
+ ${log4j.version}
+ test
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+ ${log4j.version}
+ test
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
+ test
org.mockito
@@ -275,7 +502,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -303,5 +530,4 @@
-
diff --git a/pipeline/cache-indexer/pom.xml b/pipeline/cache-indexer/pom.xml
index f5063f86..3099ac34 100644
--- a/pipeline/cache-indexer/pom.xml
+++ b/pipeline/cache-indexer/pom.xml
@@ -27,8 +27,21 @@
com.fasterxml.jackson.core
jackson-databind
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
+
org.sunbird.obsrv
framework
@@ -55,6 +68,56 @@
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.sunbird.obsrv
@@ -75,6 +138,52 @@
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -82,6 +191,29 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
com.github.codemonstur
@@ -95,6 +227,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -117,14 +255,57 @@
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
@@ -188,7 +369,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -225,7 +406,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -234,6 +415,7 @@
test
+ test
test
diff --git a/pipeline/dataset-router/pom.xml b/pipeline/dataset-router/pom.xml
index 646de8b0..0c348d64 100644
--- a/pipeline/dataset-router/pom.xml
+++ b/pipeline/dataset-router/pom.xml
@@ -32,6 +32,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -60,8 +75,28 @@
com.google.guava
guava
+
+ org.mozilla
+ rhino
+
+
+
+ com.sun.mail
+ mailapi
+
+
+
+ com.sun.mail
+ mailapi
+ 1.6.8
+
+
+ org.mozilla
+ rhino
+ 1.7.15.1
+
com.google.guava
guava
@@ -86,18 +121,120 @@
kafka-clients
${kafka.version}
test
+
+
+ org.lz4
+ lz4-java
+
+
-
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
+
+
org.apache.flink
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -105,6 +242,29 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
com.github.codemonstur
@@ -115,14 +275,57 @@
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
+ test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
org.apache.flink
@@ -130,6 +333,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -210,7 +419,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -247,7 +456,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -256,6 +465,7 @@
test
+ test
test
diff --git a/pipeline/denormalizer/pom.xml b/pipeline/denormalizer/pom.xml
index 00261ebd..3eeb3e5d 100644
--- a/pipeline/denormalizer/pom.xml
+++ b/pipeline/denormalizer/pom.xml
@@ -32,6 +32,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -53,12 +68,68 @@
kafka-clients
${kafka.version}
test
+
+
+ org.lz4
+ lz4-java
+
+
-
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.sunbird.obsrv
@@ -79,6 +150,52 @@
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -86,6 +203,29 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
com.github.codemonstur
@@ -96,14 +236,57 @@
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
org.apache.flink
@@ -111,6 +294,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -193,7 +382,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -230,7 +419,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -239,6 +428,7 @@
test
+ test
test
diff --git a/pipeline/extractor/pom.xml b/pipeline/extractor/pom.xml
index 8a8abf28..ad57b861 100644
--- a/pipeline/extractor/pom.xml
+++ b/pipeline/extractor/pom.xml
@@ -31,6 +31,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -53,12 +68,68 @@
kafka-clients
${kafka.version}
test
+
+
+ org.lz4
+ lz4-java
+
+
-
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.sunbird.obsrv
@@ -79,6 +150,52 @@
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -86,6 +203,29 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
com.github.codemonstur
@@ -96,14 +236,57 @@
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
org.apache.flink
@@ -111,6 +294,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -192,7 +381,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -229,7 +418,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -238,6 +427,7 @@
test
+ test
test
diff --git a/pipeline/hudi-connector/pom.xml b/pipeline/hudi-connector/pom.xml
deleted file mode 100644
index aaafc0cc..00000000
--- a/pipeline/hudi-connector/pom.xml
+++ /dev/null
@@ -1,327 +0,0 @@
-
-
- 4.0.0
-
- pipeline
- org.sunbird.obsrv
- 1.0
- ../pom.xml
-
- hudi-connector
- 1.0.0
- Hudi Connector
-
- UTF-8
- 1.4.0
- 2.13.4.20221013
- 1.17.2
-
-
-
-
-
- com.fasterxml.jackson
- jackson-bom
- ${jackson-bom.version}
- pom
- import
-
-
-
-
-
-
- org.apache.flink
- flink-streaming-scala_${scala.maj.version}
- ${flink.version}
- provided
-
-
- com.fasterxml.jackson.core
- jackson-databind
-
-
-
-
- com.fasterxml.jackson.core
- jackson-databind
- provided
-
-
- com.fasterxml.jackson.core
- jackson-annotations
- provided
-
-
- com.fasterxml.jackson.core
- jackson-core
- provided
-
-
- org.sunbird.obsrv
- framework
- 1.0.0
-
-
- com.fasterxml.jackson.core
- jackson-core
-
-
- com.fasterxml.jackson.core
- jackson-databind
-
-
-
-
- org.sunbird.obsrv
- dataset-registry
- 1.0.0
-
-
- org.apache.kafka
- kafka-clients
-
-
-
-
- org.apache.kafka
- kafka-clients
- ${kafka.version}
-
-
- org.apache.hudi
- hudi-flink1.17-bundle
- 1.0.2
-
-
- org.apache.hadoop
- hadoop-common
-
-
- com.fasterxml.jackson.core
- jackson-core
-
-
- com.fasterxml.jackson.core
- jackson-databind
-
-
- com.fasterxml.jackson.core
- jackson-annotations
-
-
- org.slf4j
- slf4j-log4j12
-
-
-
-
- org.apache.flink
- flink-table-api-java-bridge
- ${flink.version}
- provided
-
-
- io.github.classgraph
- classgraph
- 4.8.168
-
-
- org.apache.flink
- flink-connector-hive_${scala.maj.version}
- ${flink.version}
-
-
- org.apache.hive
- hive-metastore
- 3.1.3
-
-
- org.apache.hadoop
- hadoop-common
-
-
- com.fasterxml.jackson.core
- jackson-core
-
-
- com.fasterxml.jackson.core
- jackson-databind
-
-
- com.fasterxml.jackson.core
- jackson-annotations
-
-
-
-
- org.apache.hive
- hive-exec
- 3.1.3
-
-
- com.fasterxml.jackson.core
- jackson-core
-
-
- com.fasterxml.jackson.core
- jackson-databind
-
-
- com.fasterxml.jackson.core
- jackson-annotations
-
-
- org.apache.hadoop
- hadoop-common
-
-
- org.apache.logging.log4j
- log4j-slf4j-impl
-
-
-
-
- org.apache.flink
- flink-statebackend-rocksdb
- 1.17.2
- provided
-
-
- org.scalatest
- scalatest_2.12
- 3.2.17
- test
-
-
- org.sunbird.obsrv
- dataset-registry
- 1.0.0
-
-
-
-
- src/main/scala
-
-
-
- org.apache.maven.plugins
- maven-compiler-plugin
- 3.8.1
-
- 11
-
-
-
- org.apache.maven.plugins
- maven-shade-plugin
- 3.2.1
-
-
-
- package
-
- shade
-
-
- false
-
-
- com.google.code.findbugs:jsr305
-
-
-
-
-
- *:*
-
- META-INF/*.SF
- META-INF/*.DSA
- META-INF/*.RSA
- core-site.xml
-
-
-
-
-
- org.sunbird.obsrv.streaming.HudiConnectorStreamTask
-
-
-
- reference.conf
-
-
-
-
-
-
-
-
- net.alchim31.maven
- scala-maven-plugin
- 4.4.0
-
- ${java.target.runtime}
- ${java.target.runtime}
- ${scala.version}
- false
-
-
-
- scala-compile-first
- process-resources
-
- add-source
- compile
-
-
-
- scala-test-compile
- process-test-resources
-
- testCompile
-
-
-
-
-
-
- maven-surefire-plugin
- 2.22.2
-
- true
-
-
-
-
- org.scalatest
- scalatest-maven-plugin
- 1.0
-
- ${project.build.directory}/surefire-reports
- .
- hudi-connector-testsuite.txt
-
-
-
- test
-
- test
-
-
-
-
-
- org.scoverage
- scoverage-maven-plugin
- ${scoverage.plugin.version}
-
- ${scala.version}
- true
- true
-
-
-
-
-
diff --git a/pipeline/hudi-connector/src/main/resources/core-site.xml b/pipeline/hudi-connector/src/main/resources/core-site.xml
deleted file mode 100644
index c15df562..00000000
--- a/pipeline/hudi-connector/src/main/resources/core-site.xml
+++ /dev/null
@@ -1,33 +0,0 @@
-
-
-
-
-
- fs.s3a.impl
- org.apache.hadoop.fs.s3a.S3AFileSystem
-
-
- fs.s3a.endpoint
- http://localhost:4566
-
-
- fs.s3a.access.key
- test
-
-
- fs.s3a.secret.key
- testSecret
-
-
- fs.s3a.path.style.access
- true
-
-
- fs.s3a.connection.ssl.enabled
- false
-
-
-
-
-
-
\ No newline at end of file
diff --git a/pipeline/hudi-connector/src/main/resources/hudi-writer.conf b/pipeline/hudi-connector/src/main/resources/hudi-writer.conf
deleted file mode 100644
index 093bb2a4..00000000
--- a/pipeline/hudi-connector/src/main/resources/hudi-writer.conf
+++ /dev/null
@@ -1,52 +0,0 @@
-include "baseconfig.conf"
-
-kafka {
- input.topic = "hudi.connector.in"
- output.topic = "hudi.connector.out"
- output.invalid.topic = "failed"
- event.max.size = "1048576" # Max is only 1MB
- groupId = "hudi-writer-group"
- producer {
- max-request-size = 5242880
- }
-}
-
-task {
- checkpointing.compressed = true
- checkpointing.interval = 30000
- checkpointing.pause.between.seconds = 30000
- restart-strategy.attempts = 3
- restart-strategy.delay = 30000 # in milli-seconds
- parallelism = 1
- consumer.parallelism = 1
- downstream.operators.parallelism = 1
-}
-
-hudi {
- hms {
- enabled = true
- uri = "thrift://localhost:9083"
- database {
- name = "obsrv"
- username = "postgres"
- password = "postgres"
- }
- }
- table {
- type = "MERGE_ON_READ"
- base.path = "s3a://obsrv"
- }
- compaction.enabled = true
- metadata.enabled = true
- write {
- tasks = 1
- task.max.memory = 512
- compaction.max.memory = 100
- }
- write.batch.size = 16
- compaction.tasks = 1
- delta.commits = 2
- delta.seconds = 600
- compression.codec = "snappy"
- index.type = "BLOOM"
-}
\ No newline at end of file
diff --git a/pipeline/hudi-connector/src/main/resources/schemas/schema.json b/pipeline/hudi-connector/src/main/resources/schemas/schema.json
deleted file mode 100644
index 177c957a..00000000
--- a/pipeline/hudi-connector/src/main/resources/schemas/schema.json
+++ /dev/null
@@ -1,108 +0,0 @@
-{
- "dataset": "financial_transactions",
- "schema": {
- "table": "financial_transactions",
- "partitionColumn": "receiver_ifsc_code",
- "timestampColumn": "txn_date",
- "primaryKey": "txn_id",
- "columnSpec": [
- {
- "name": "receiver_account_number",
- "type": "string"
- },
- {
- "name": "receiver_ifsc_code",
- "type": "string"
- },
- {
- "name": "sender_account_number",
- "type": "string"
- },
- {
- "name": "sender_contact_email",
- "type": "string"
- },
- {
- "name": "sender_ifsc_code",
- "type": "string"
- },
- {
- "name": "currency",
- "type": "string"
- },
- {
- "name": "txn_amount",
- "type": "int"
- },
- {
- "name": "txn_date",
- "type": "string"
- },
- {
- "name": "txn_id",
- "type": "string"
- },
- {
- "name": "txn_status",
- "type": "string"
- },
- {
- "name": "txn_type",
- "type": "string"
- }
- ]
- },
- "inputFormat": {
- "type": "json",
- "flattenSpec": {
- "fields": [
- {
- "type": "root",
- "name": "receiver_account_number"
- },
- {
- "type": "path",
- "name": "sender_account_number",
- "expr": "$.sender.account_number"
- },
- {
- "type": "path",
- "name": "sender_ifsc_code",
- "expr": "$.sender.ifsc_code"
- },
- {
- "type": "root",
- "name": "receiver_ifsc_code"
- },
- {
- "type": "root",
- "name": "sender_contact_email"
- },
- {
- "type": "root",
- "name": "currency"
- },
- {
- "type": "root",
- "name": "txn_amount"
- },
- {
- "type": "root",
- "name": "txn_date"
- },
- {
- "type": "root",
- "name": "txn_id"
- },
- {
- "type": "root",
- "name": "txn_status"
- },
- {
- "type": "root",
- "name": "txn_type"
- }
- ]
- }
- }
-}
\ No newline at end of file
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/function/RowDataConverterFunction.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/function/RowDataConverterFunction.scala
deleted file mode 100644
index f85979e6..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/function/RowDataConverterFunction.scala
+++ /dev/null
@@ -1,72 +0,0 @@
-package org.sunbird.obsrv.function
-
-import org.apache.flink.api.common.functions.RichMapFunction
-import org.apache.flink.configuration.Configuration
-import org.apache.flink.formats.common.TimestampFormat
-import org.apache.flink.formats.json.JsonToRowDataConverters
-import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper
-import org.apache.flink.table.data.RowData
-import org.slf4j.LoggerFactory
-import org.sunbird.obsrv.core.util.JSONUtil
-import org.sunbird.obsrv.streaming.HudiConnectorConfig
-import org.sunbird.obsrv.util.{HMetrics, HudiSchemaParser, ScalaGauge}
-
-import scala.collection.mutable.{Map => MMap}
-
-class RowDataConverterFunction(config: HudiConnectorConfig, datasetId: String)
- extends RichMapFunction[MMap[String, AnyRef], RowData] {
-
- private val logger = LoggerFactory.getLogger(classOf[RowDataConverterFunction])
-
- private var metrics: HMetrics = _
- private var jsonToRowDataConverters: JsonToRowDataConverters = _
- private var objectMapper: ObjectMapper = _
- private var hudiSchemaParser: HudiSchemaParser = _
-
- override def open(parameters: Configuration): Unit = {
- super.open(parameters)
-
- metrics = new HMetrics()
- jsonToRowDataConverters = new JsonToRowDataConverters(false, true, TimestampFormat.SQL)
- objectMapper = new ObjectMapper()
- hudiSchemaParser = new HudiSchemaParser()
-
- getRuntimeContext.getMetricGroup
- .addGroup(config.jobName)
- .addGroup(datasetId)
- .gauge[Long, ScalaGauge[Long]](config.inputEventCountMetric, ScalaGauge[Long](() =>
- metrics.getAndReset(datasetId, config.inputEventCountMetric)
- ))
-
- getRuntimeContext.getMetricGroup
- .addGroup(config.jobName)
- .addGroup(datasetId)
- .gauge[Long, ScalaGauge[Long]](config.failedEventCountMetric, ScalaGauge[Long](() =>
- metrics.getAndReset(datasetId, config.failedEventCountMetric)
- ))
- }
-
- override def map(event: MMap[String, AnyRef]): RowData = {
- try {
- if (event.nonEmpty) {
- metrics.increment(datasetId, config.inputEventCountMetric, 1)
- }
- val rowData = convertToRowData(event)
- rowData
- } catch {
- case ex: Exception =>
- metrics.increment(datasetId, config.failedEventCountMetric, 1)
- logger.error("Failed to process record", ex)
- throw ex
- }
- }
-
- def convertToRowData(data: MMap[String, AnyRef]): RowData = {
- val eventJson = JSONUtil.serialize(data)
- val flattenedData = hudiSchemaParser.parseJson(datasetId, eventJson)
- val rowType = hudiSchemaParser.rowTypeMap(datasetId)
- val converter: JsonToRowDataConverters.JsonToRowDataConverter =
- jsonToRowDataConverters.createRowConverter(rowType)
- converter.convert(objectMapper.readTree(JSONUtil.serialize(flattenedData))).asInstanceOf[RowData]
- }
-}
\ No newline at end of file
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorConfig.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorConfig.scala
deleted file mode 100644
index 1bf12167..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorConfig.scala
+++ /dev/null
@@ -1,75 +0,0 @@
-package org.sunbird.obsrv.streaming
-
-import com.typesafe.config.Config
-import org.apache.flink.api.common.typeinfo.TypeInformation
-import org.apache.flink.api.java.typeutils.TypeExtractor
-import org.apache.flink.streaming.api.scala.OutputTag
-import org.apache.hudi.common.model.HoodieTableType
-import org.sunbird.obsrv.core.streaming.BaseJobConfig
-
-import scala.collection.mutable
-
-class HudiConnectorConfig(override val config: Config) extends BaseJobConfig[mutable.Map[String, AnyRef]](config, "Flink-Hudi-Connector") {
-
- implicit val mapTypeInfo: TypeInformation[mutable.Map[String, AnyRef]] = TypeExtractor.getForClass(classOf[mutable.Map[String, AnyRef]])
-
- override def inputTopic(): String = config.getString("kafka.input.topic")
-
- val kafkaDefaultOutputTopic: String = config.getString("kafka.output.topic")
-
- override def inputConsumer(): String = config.getString("kafka.groupId")
-
- override def successTag(): OutputTag[mutable.Map[String, AnyRef]] = OutputTag[mutable.Map[String, AnyRef]]("dummy-events")
-
- override def failedEventsOutputTag(): OutputTag[mutable.Map[String, AnyRef]] = OutputTag[mutable.Map[String, AnyRef]]("failed-events")
-
- val kafkaInvalidTopic: String = config.getString("kafka.output.invalid.topic")
-
- val invalidEventsOutputTag: OutputTag[mutable.Map[String, AnyRef]] = OutputTag[mutable.Map[String, AnyRef]]("invalid-events")
- val validEventsOutputTag: OutputTag[mutable.Map[String, AnyRef]] = OutputTag[mutable.Map[String, AnyRef]]("valid-events")
-
- val invalidEventProducer = "invalid-events-sink"
-
-
- val hudiTableType: String =
- if (config.getString("hudi.table.type").equalsIgnoreCase("MERGE_ON_READ"))
- HoodieTableType.MERGE_ON_READ.name()
- else if (config.getString("hudi.table.type").equalsIgnoreCase("COPY_ON_WRITE"))
- HoodieTableType.COPY_ON_WRITE.name()
- else HoodieTableType.MERGE_ON_READ.name()
-
- val hudiBasePath: String = config.getString("hudi.table.base.path")
-
- val hmsEnabled: Boolean = if (config.hasPath("hudi.hms.enabled")) config.getBoolean("hudi.hms.enabled") else false
- val hmsUsername: String = config.getString("hudi.hms.database.username")
- val hmsPassword: String = config.getString("hudi.hms.database.password")
- val hmsDatabaseName: String = config.getString("hudi.hms.database.name")
- val hmsURI: String = config.getString("hudi.hms.uri")
-
- val hudiWriteTasks: Int = config.getInt("hudi.write.tasks")
- val hudiCompactionTasks: Int = config.getInt("hudi.compaction.tasks")
- val hudiWriteBatchSize: Int = config.getInt("hudi.write.batch.size")
- val deltaCommits: Int = config.getInt("hudi.delta.commits")
- val compactionDeltaSeconds: Int = config.getInt("hudi.delta.seconds")
- val compressionCodec: String = config.getString("hudi.compression.codec")
- val hudiCompactionEnabled: Boolean = config.getBoolean("hudi.compaction.enabled")
- val hudiMetadataEnabled: Boolean = config.getBoolean("hudi.metadata.enabled")
- val hudiIndexType: String = config.getString("hudi.index.type")
-
- // Memory
- val hudiWriteTaskMemory: Int = config.getInt("hudi.write.task.max.memory")
- val hudiCompactionTaskMemory: Int = config.getInt("hudi.write.compaction.max.memory")
- val hudiFsAtomicCreationSupport: String = config.getString("hudi.fs.atomic_creation.support")
-
- // Metrics
-
- val inputEventCountMetric = "input-event-count"
- val failedEventCountMetric = "failed-event-count"
-
- // Metrics Exporter
- val metricsReportType: String = config.getString("metrics.reporter.type")
- val metricsReporterHost: String = config.getString("metrics.reporter.host")
- val metricsReporterPort: String = config.getString("metrics.reporter.port")
-
-
-}
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorStreamTask.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorStreamTask.scala
deleted file mode 100644
index cc28214c..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/HudiConnectorStreamTask.scala
+++ /dev/null
@@ -1,176 +0,0 @@
-package org.sunbird.obsrv.streaming
-
-import com.typesafe.config.ConfigFactory
-import org.apache.commons.lang3.StringUtils
-import org.apache.flink.api.common.typeinfo.TypeInformation
-import org.apache.flink.api.java.typeutils.TypeExtractor
-import org.apache.flink.api.java.utils.ParameterTool
-import org.apache.flink.configuration.Configuration
-import org.apache.flink.contrib.streaming.state.EmbeddedRocksDBStateBackend
-import org.apache.flink.streaming.api.datastream.{DataStream, DataStreamSink}
-import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment
-import org.apache.hudi.common.config.TimestampKeyGeneratorConfig
-import org.apache.hudi.configuration.{FlinkOptions, OptionsResolver}
-import org.apache.hudi.sink.utils.Pipelines
-import org.apache.hudi.util.AvroSchemaConverter
-import org.slf4j.LoggerFactory
-import org.sunbird.obsrv.core.model.Constants
-import org.sunbird.obsrv.core.streaming.{BaseStreamTask, FlinkKafkaConnector}
-import org.sunbird.obsrv.core.util.FlinkUtil
-import org.sunbird.obsrv.function.RowDataConverterFunction
-import org.sunbird.obsrv.registry.DatasetRegistry
-import org.sunbird.obsrv.util.HudiSchemaParser
-import org.apache.hudi.config.HoodieWriteConfig.SCHEMA_ALLOW_AUTO_EVOLUTION_COLUMN_DROP
-import org.apache.hudi.common.config.HoodieCommonConfig.SCHEMA_EVOLUTION_ENABLE
-import org.apache.hudi.common.table.HoodieTableConfig.DROP_PARTITION_COLUMNS
-
-import java.io.File
-import java.sql.Timestamp
-import java.time.LocalDateTime
-import java.time.format.DateTimeFormatter
-import scala.collection.mutable
-import scala.collection.mutable.{Map => MMap}
-
-class HudiConnectorStreamTask(config: HudiConnectorConfig, kafkaConnector: FlinkKafkaConnector) extends BaseStreamTask[mutable.Map[String, AnyRef]] {
-
- implicit val mutableMapTypeInfo: TypeInformation[MMap[String, AnyRef]] = TypeExtractor.getForClass(classOf[MMap[String, AnyRef]])
- private val logger = LoggerFactory.getLogger(classOf[HudiConnectorStreamTask])
- def process(): Unit = {
- implicit val env: StreamExecutionEnvironment = FlinkUtil.getExecutionContext(config)
- env.setStateBackend(new EmbeddedRocksDBStateBackend)
- process(env)
- }
-
- override def processStream(dataStream: DataStream[mutable.Map[String, AnyRef]]): DataStream[mutable.Map[String, AnyRef]] = {
- null
- }
-
- def process(env: StreamExecutionEnvironment): Unit = {
- val schemaParser = new HudiSchemaParser()
- val dataSourceConfig = DatasetRegistry.getAllDatasources().filter(f => f.`type`.nonEmpty && f.`type`.equalsIgnoreCase(Constants.DATALAKE_TYPE) && f.status.equalsIgnoreCase("Live"))
- dataSourceConfig.map{ dataSource =>
- val datasetId = dataSource.datasetId
- val dataStream = getMapDataStream(env, config, List(datasetId), config.kafkaConsumerProperties(), consumerSourceName = s"kafka-${datasetId}", kafkaConnector)
- .map(new RowDataConverterFunction(config, datasetId))
- .setParallelism(config.downstreamOperatorsParallelism)
-
- val conf: Configuration = new Configuration()
- setHudiBaseConfigurations(conf)
- setDatasetConf(conf, datasetId, schemaParser)
- logger.info("conf: " + conf.toMap.toString)
- val rowType = schemaParser.rowTypeMap(datasetId)
-
- val hoodieRecordDataStream = Pipelines.bootstrap(conf, rowType, dataStream)
- val pipeline = Pipelines.hoodieStreamWrite(conf, rowType, hoodieRecordDataStream)
- if (OptionsResolver.needsAsyncCompaction(conf)) {
- Pipelines.compact(conf, pipeline).setParallelism(config.downstreamOperatorsParallelism)
- } else {
- Pipelines.clean(conf, pipeline).setParallelism(config.downstreamOperatorsParallelism)
- }
- }.orElse(List(addDefaultOperator(env, config, kafkaConnector)))
- env.execute("Flink-Hudi-Connector")
- }
-
- def addDefaultOperator(env: StreamExecutionEnvironment, config: HudiConnectorConfig, kafkaConnector: FlinkKafkaConnector): DataStreamSink[mutable.Map[String, AnyRef]] = {
- val dataStreamSink: DataStreamSink[mutable.Map[String, AnyRef]] = getMapDataStream(env, config, kafkaConnector)
- .sinkTo(kafkaConnector.kafkaSink[mutable.Map[String, AnyRef]](config.kafkaDefaultOutputTopic))
- .name(s"hudi-connector-default-sink").uid(s"hudi-connector-default-sink")
- .setParallelism(config.downstreamOperatorsParallelism)
- dataStreamSink
- }
-
- def setDatasetConf(conf: Configuration, dataset: String, schemaParser: HudiSchemaParser): Unit = {
- val datasetSchema = schemaParser.hudiSchemaMap(dataset)
- val rowType = schemaParser.rowTypeMap(dataset)
- val avroSchema = AvroSchemaConverter.convertToSchema(rowType, dataset.replace("-", "_"))
- conf.setString(FlinkOptions.PATH.key, s"${config.hudiBasePath}/${datasetSchema.schema.table}")
- conf.setString("hoodie.base.path", s"${config.hudiBasePath}/${datasetSchema.schema.table}")
- conf.setString(FlinkOptions.TABLE_NAME, datasetSchema.schema.table)
- conf.setString(FlinkOptions.RECORD_KEY_FIELD.key, datasetSchema.schema.primaryKey)
- conf.setString(FlinkOptions.PRECOMBINE_FIELD.key, datasetSchema.schema.timestampColumn)
- conf.setString(FlinkOptions.PARTITION_PATH_FIELD.key, datasetSchema.schema.partitionColumn)
- conf.setString(FlinkOptions.SOURCE_AVRO_SCHEMA.key, avroSchema.toString)
- conf.setBoolean("hoodie.metrics.on", true)
- if (config.metricsReportType.equalsIgnoreCase("PROMETHEUS_PUSHGATEWAY")) {
- conf.setString("hoodie.metrics.reporter.type", config.metricsReportType)
- conf.setString("hoodie.metrics.pushgateway.host", config.metricsReporterHost)
- conf.setString("hoodie.metrics.pushgateway.port", config.metricsReporterPort)
- }
- if (config.metricsReportType.equalsIgnoreCase("JMX")) {
- conf.setString("hoodie.metrics.reporter.type", config.metricsReportType)
- conf.setString("hoodie.metrics.jmx.host", config.metricsReporterHost)
- conf.setString("hoodie.metrics.jmx.port", config.metricsReporterPort)
- }
- val partitionField = datasetSchema.schema.columnSpec.filter(f => f.name.equalsIgnoreCase(datasetSchema.schema.partitionColumn)).head
- if(partitionField.`type`.equalsIgnoreCase("timestamp") || partitionField.`type`.equalsIgnoreCase("epoch")) {
- conf.setString(FlinkOptions.PARTITION_PATH_FIELD.key, datasetSchema.schema.partitionColumn + "_partition")
- }
-
- if (config.hmsEnabled) {
- conf.setString("hive_sync.table", datasetSchema.schema.table)
- }
- }
-
- private def setHudiBaseConfigurations(conf: Configuration): Unit = {
- conf.setString(FlinkOptions.TABLE_TYPE.key, config.hudiTableType)
- conf.setBoolean(FlinkOptions.METADATA_ENABLED.key, config.hudiMetadataEnabled)
-
- conf.setDouble(FlinkOptions.WRITE_BATCH_SIZE.key, config.hudiWriteBatchSize)
- conf.setInteger(FlinkOptions.COMPACTION_TASKS, config.downstreamOperatorsParallelism)
- conf.setBoolean(FlinkOptions.COMPACTION_SCHEDULE_ENABLED.key, config.hudiCompactionEnabled)
- conf.setInteger(FlinkOptions.WRITE_TASKS, config.downstreamOperatorsParallelism)
- conf.setInteger(FlinkOptions.COMPACTION_DELTA_COMMITS, config.deltaCommits)
- conf.setString(FlinkOptions.COMPACTION_TRIGGER_STRATEGY, FlinkOptions.NUM_OR_TIME)
- conf.setInteger(FlinkOptions.COMPACTION_DELTA_SECONDS, config.compactionDeltaSeconds)
- conf.setBoolean(FlinkOptions.COMPACTION_ASYNC_ENABLED, true)
-
- conf.setString("hoodie.fs.atomic_creation.support", config.hudiFsAtomicCreationSupport)
- conf.setString(FlinkOptions.HIVE_SYNC_TABLE_PROPERTIES, "hoodie.datasource.write.drop.partition.columns=true")
- conf.setBoolean(DROP_PARTITION_COLUMNS.key, true)
- conf.setBoolean(SCHEMA_ALLOW_AUTO_EVOLUTION_COLUMN_DROP.key(), true); // Enable dropping columns
- conf.setBoolean(SCHEMA_EVOLUTION_ENABLE.key(), true); // Enable schema evolution
- conf.setString(FlinkOptions.PAYLOAD_CLASS_NAME, "org.apache.hudi.common.model.PartialUpdateAvroPayload")
- conf.setString("hoodie.parquet.compression.codec", config.compressionCodec)
-
- // Index Type Configurations
- conf.setString(FlinkOptions.INDEX_TYPE, config.hudiIndexType)
- conf.setInteger(FlinkOptions.BUCKET_ASSIGN_TASKS, config.downstreamOperatorsParallelism)
- conf.setDouble(FlinkOptions.WRITE_TASK_MAX_SIZE, config.hudiWriteTaskMemory)
- conf.setInteger(FlinkOptions.COMPACTION_MAX_MEMORY, config.hudiCompactionTaskMemory)
-
- if (config.hmsEnabled) {
- conf.setBoolean("hive_sync.enabled", config.hmsEnabled)
- conf.setString(FlinkOptions.HIVE_SYNC_DB.key(), config.hmsDatabaseName)
- conf.setString("hive_sync.username", config.hmsUsername)
- conf.setString("hive_sync.password", config.hmsPassword)
- conf.setString("hive_sync.mode", "hms")
- conf.setBoolean("hive_sync.use_jdbc", false)
- conf.setString(FlinkOptions.HIVE_SYNC_METASTORE_URIS.key(), config.hmsURI)
- conf.setString("hoodie.fs.atomic_creation.support", config.hudiFsAtomicCreationSupport)
- conf.setBoolean(FlinkOptions.HIVE_SYNC_SUPPORT_TIMESTAMP, true)
- }
-
- }
-
-}
-
-object HudiConnectorStreamTask {
- def main(args: Array[String]): Unit = {
- val configFilePath = Option(ParameterTool.fromArgs(args).get("config.file.path"))
- val config = configFilePath.map {
- path => ConfigFactory.parseFile(new File(path)).resolve()
- }.getOrElse(ConfigFactory.load("hudi-writer.conf").withFallback(ConfigFactory.systemEnvironment()))
- val hudiWriterConfig = new HudiConnectorConfig(config)
- val kafkaUtil = new FlinkKafkaConnector(hudiWriterConfig)
- val task = new HudiConnectorStreamTask(hudiWriterConfig, kafkaUtil)
- task.process()
- }
-
- def getTimestamp(ts: String): Timestamp = {
- val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSXXX")
- val localDateTime = if (StringUtils.isNotBlank(ts))
- LocalDateTime.from(formatter.parse(ts))
- else LocalDateTime.now
- Timestamp.valueOf(localDateTime)
- }
-}
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/TestTimestamp.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/TestTimestamp.scala
deleted file mode 100644
index 1c9876c0..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/streaming/TestTimestamp.scala
+++ /dev/null
@@ -1,19 +0,0 @@
-package org.sunbird.obsrv.streaming
-
-import java.sql.Timestamp
-import java.time.{LocalDateTime, ZoneOffset}
-import java.time.format.DateTimeFormatter
-
-object TestTimestamp {
-
- def main(args: Array[String]): Unit = {
- val timestampAsString = "2023-10-15T03:56:27.522+05:30"
- val pattern = "yyyy-MM-dd'T'hh:mm:ss.SSSZ"
- val formatter = DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSSXXX")
- val localDateTime = LocalDateTime.from(formatter.parse(timestampAsString))
- val timestamp = Timestamp.valueOf(localDateTime)
- println("Timestamp: " + timestamp.toString)
-
- }
-
-}
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HMetrics.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HMetrics.scala
deleted file mode 100644
index 3f3d1b1c..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HMetrics.scala
+++ /dev/null
@@ -1,31 +0,0 @@
-package org.sunbird.obsrv.util
-
-
-import org.apache.flink.metrics.Gauge
-
-import scala.collection.concurrent.TrieMap
-
-class HMetrics {
- private val metricStore = TrieMap[(String, String), Long]()
-
- def increment(dataset: String, metric: String, value: Long): Unit = {
- metricStore.synchronized {
- val key = (dataset, metric)
- val current = metricStore.getOrElse(key, 0L)
- metricStore.put(key, current + value)
- }
- }
-
- def getAndReset(dataset: String, metric: String): Long = {
- metricStore.synchronized {
- val key = (dataset, metric)
- val current = metricStore.getOrElse(key, 0L)
- metricStore.remove(key)
- current
- }
- }
-}
-
-case class ScalaGauge[T](getValueFn: () => T) extends Gauge[T] {
- override def getValue: T = getValueFn()
-}
diff --git a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HudiSchemaParser.scala b/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HudiSchemaParser.scala
deleted file mode 100644
index 5161a84b..00000000
--- a/pipeline/hudi-connector/src/main/scala/org/sunbird/obsrv/util/HudiSchemaParser.scala
+++ /dev/null
@@ -1,140 +0,0 @@
-package org.sunbird.obsrv.util
-
-import com.fasterxml.jackson.annotation.JsonInclude.Include
-import com.fasterxml.jackson.core.JsonGenerator.Feature
-import com.fasterxml.jackson.databind.json.JsonMapper
-import com.fasterxml.jackson.databind.{DeserializationFeature, JsonNode, ObjectMapper, SerializationFeature}
-import com.fasterxml.jackson.module.scala.DefaultScalaModule
-import org.apache.flink.table.types.logical.{BigIntType, BooleanType, DoubleType, IntType, LogicalType, MapType, RowType, VarCharType, TimestampType, DateType}
-import org.slf4j.LoggerFactory
-import org.sunbird.obsrv.core.model.Constants
-import org.sunbird.obsrv.core.util.JSONUtil
-import org.sunbird.obsrv.registry.DatasetRegistry
-import java.sql.Timestamp
-import java.text.SimpleDateFormat
-import java.util.Date
-import scala.collection.mutable
-
-
-case class HudiSchemaSpec(dataset: String, schema: Schema, inputFormat: InputFormat)
-case class Schema(table: String, partitionColumn: String, timestampColumn: String, primaryKey: String, columnSpec: List[ColumnSpec])
-case class ColumnSpec(name: String, `type`: String)
-case class InputFormat(`type`: String, flattenSpec: Option[JsonFlattenSpec] = None, columns: Option[List[String]] = None)
-case class JsonFlattenSpec(fields: List[JsonFieldParserSpec])
-case class JsonFieldParserSpec(`type`: String, name: String, expr: Option[String] = None)
-
-class HudiSchemaParser {
-
- private val logger = LoggerFactory.getLogger(classOf[HudiSchemaParser])
-
- @transient private val objectMapper = JsonMapper.builder()
- .addModule(DefaultScalaModule)
- .disable(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES)
- .disable(SerializationFeature.FAIL_ON_EMPTY_BEANS)
- .enable(Feature.WRITE_BIGDECIMAL_AS_PLAIN)
- .build()
-
- val df = new SimpleDateFormat("yyyy-MM-dd")
- objectMapper.setSerializationInclusion(Include.NON_ABSENT)
-
- val hudiSchemaMap = new mutable.HashMap[String, HudiSchemaSpec]()
- val rowTypeMap = new mutable.HashMap[String, RowType]()
-
- readSchema()
-
- def readSchema(): Unit = {
- val datasourceConfig = DatasetRegistry.getAllDatasources().filter(f => f.`type`.nonEmpty && f.`type`.equalsIgnoreCase(Constants.DATALAKE_TYPE) && f.status.equalsIgnoreCase("Live"))
- datasourceConfig.map{f =>
- val hudiSchemaSpec = JSONUtil.deserialize[HudiSchemaSpec](f.ingestionSpec)
- val dataset = hudiSchemaSpec.dataset
- hudiSchemaMap.put(dataset, hudiSchemaSpec)
- rowTypeMap.put(dataset, createRowType(hudiSchemaSpec))
- }
- }
-
- private def createRowType(schema: HudiSchemaSpec): RowType = {
- val columnSpec = schema.schema.columnSpec
- val primaryKey = schema.schema.primaryKey
- val partitionColumn = schema.schema.partitionColumn
- val timeStampColumn = schema.schema.timestampColumn
- val partitionField = schema.schema.columnSpec.filter(f => f.name.equalsIgnoreCase(schema.schema.partitionColumn)).head
- val rowTypeMap = mutable.SortedMap[String, LogicalType]()
- columnSpec.sortBy(_.name).map {
- spec =>
- val isNullable = if (spec.name.matches(s"$primaryKey|$partitionColumn|$timeStampColumn")) false else true
- val columnType = spec.`type` match {
- case "string" => new VarCharType(isNullable, 20)
- case "double" => new DoubleType(isNullable)
- case "long" => new BigIntType(isNullable)
- case "int" => new IntType(isNullable)
- case "boolean" => new BooleanType(true)
- case "map[string, string]" => new MapType(new VarCharType(), new VarCharType())
- case "epoch" => new BigIntType(isNullable)
- case _ => new VarCharType(isNullable, 20)
- }
- rowTypeMap.put(spec.name, columnType)
- }
- if(partitionField.`type`.equalsIgnoreCase("timestamp") || partitionField.`type`.equalsIgnoreCase("epoch")) {
- rowTypeMap.put(partitionField.name + "_partition", new VarCharType(false, 20))
- }
- val rowType: RowType = RowType.of(false, rowTypeMap.values.toArray, rowTypeMap.keySet.toArray)
- logger.info("rowType: " + rowType)
- rowType
- }
-
- def parseJson(dataset: String, event: String): mutable.Map[String, Any] = {
- val parserSpec = hudiSchemaMap.get(dataset)
- val jsonNode = objectMapper.readTree(event)
- val flattenedEventData = mutable.Map[String, Any]()
- parserSpec.map { spec =>
- val columnSpec = spec.schema.columnSpec
- val partitionField = spec.schema.columnSpec.filter(f => f.name.equalsIgnoreCase(spec.schema.partitionColumn)).head
- spec.inputFormat.flattenSpec.map {
- flattenSpec =>
- flattenSpec.fields.map {
- field =>
- val node = retrieveFieldFromJson(jsonNode, field)
- node.map {
- nodeValue =>
- try {
- val fieldDataType = columnSpec.filter(_.name.equalsIgnoreCase(field.name)).head.`type`
- val fieldValue = fieldDataType match {
- case "string" => objectMapper.treeToValue(nodeValue, classOf[String])
- case "int" => objectMapper.treeToValue(nodeValue, classOf[Int])
- case "long" => objectMapper.treeToValue(nodeValue, classOf[Long])
- case "double" => objectMapper.treeToValue(nodeValue, classOf[Double])
- case "epoch" => objectMapper.treeToValue(nodeValue, classOf[Long])
- case _ => objectMapper.treeToValue(nodeValue, classOf[String])
- }
- if(field.name.equalsIgnoreCase(partitionField.name)){
- if(fieldDataType.equalsIgnoreCase("timestamp")) {
- flattenedEventData.put(field.name + "_partition", df.format(objectMapper.treeToValue(nodeValue, classOf[Timestamp])))
- }
- else if(fieldDataType.equalsIgnoreCase("epoch")) {
- flattenedEventData.put(field.name + "_partition", df.format(objectMapper.treeToValue(nodeValue, classOf[Long])))
- }
- }
- flattenedEventData.put(field.name, fieldValue)
- }
- catch {
- case ex: Exception =>
- // logger.debug("Hudi Schema Parser - Exception: ", ex.getMessage)
- flattenedEventData.put(field.name, null)
- }
-
- }.orElse(flattenedEventData.put(field.name, null))
- }
- }
- }
- // logger.debug("flattenedEventData: " + flattenedEventData)
- flattenedEventData
- }
-
- def retrieveFieldFromJson(jsonNode: JsonNode, field: JsonFieldParserSpec): Option[JsonNode] = {
- if (field.`type`.equalsIgnoreCase("path")) {
- field.expr.map{ f => jsonNode.at(s"/${f.split("\\.").tail.mkString("/")}") }
- } else {
- Option(jsonNode.get(field.name))
- }
- }
-}
diff --git a/pipeline/pom.xml b/pipeline/pom.xml
index e3bdc523..516ce9ee 100644
--- a/pipeline/pom.xml
+++ b/pipeline/pom.xml
@@ -25,7 +25,6 @@
transformer
dataset-router
unified-pipeline
- hudi-connector
cache-indexer
preprocessor
@@ -77,7 +76,35 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
+
org.apache.flink
flink-connector-base
@@ -97,7 +124,43 @@
org.apache.flink
flink-connector-kafka
3.3.0-1.20
+
+
+ org.apache.kafka
+ kafka-clients
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+ org.apache.kafka
+ kafka-clients
+ ${kafka.version}
+
+
+ org.lz4
+ lz4-java
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+
+ org.xerial.snappy
+ snappy-java
+ 1.1.10.5
+
diff --git a/pipeline/preprocessor/pom.xml b/pipeline/preprocessor/pom.xml
index 32a493a8..df3eb67b 100644
--- a/pipeline/preprocessor/pom.xml
+++ b/pipeline/preprocessor/pom.xml
@@ -30,6 +30,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -73,7 +88,7 @@
io.github.classgraph
classgraph
- 4.8.90
+ 4.8.184
org.sunbird.obsrv
@@ -93,6 +108,52 @@
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -100,18 +161,84 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
org.sunbird.obsrv
@@ -131,12 +258,68 @@
kafka-clients
${kafka.version}
test
+
+
+ org.lz4
+ lz4-java
+
+
-
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.apache.flink
@@ -144,6 +327,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -220,7 +409,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -257,7 +446,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -266,6 +455,7 @@
test
+ test
test
diff --git a/pipeline/transformer/pom.xml b/pipeline/transformer/pom.xml
index 664d1333..f807984c 100644
--- a/pipeline/transformer/pom.xml
+++ b/pipeline/transformer/pom.xml
@@ -31,6 +31,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -66,18 +81,120 @@
kafka-clients
${kafka.version}
test
+
+
+ org.lz4
+ lz4-java
+
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.apache.flink
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -85,18 +202,84 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
com.github.codemonstur
@@ -110,6 +293,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -191,7 +380,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -228,7 +417,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -237,6 +426,7 @@
test
+ test
test
diff --git a/pipeline/unified-pipeline/pom.xml b/pipeline/unified-pipeline/pom.xml
index d44da59c..43a2403b 100644
--- a/pipeline/unified-pipeline/pom.xml
+++ b/pipeline/unified-pipeline/pom.xml
@@ -35,6 +35,21 @@
flink-streaming-scala_${scala.maj.version}
${flink.version}
provided
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ com.esotericsoftware
+ kryo
+ 4.0.3
org.sunbird.obsrv
@@ -71,11 +86,61 @@
dataset-router
1.0.0
-
+
org.apache.kafka
kafka_${scala.maj.version}
${kafka.version}
test
+
+
+ commons-beanutils
+ commons-beanutils
+
+
+ io.netty
+ netty-handler
+
+
+ io.netty
+ netty-transport-native-epoll
+
+
+ org.bitbucket.b_c
+ jose4j
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.bitbucket.b_c
+ jose4j
+ 0.9.6
+ test
+
+
+ io.netty
+ netty-transport-native-epoll
+ 4.1.130.Final
+ test
+
+
+ io.netty
+ netty-handler
+ 4.1.130.Final
+ test
+
+
+ commons-beanutils
+ commons-beanutils
+ 1.11.0
+ test
org.sunbird.obsrv
@@ -96,6 +161,52 @@
flink-test-utils
${flink.version}
test
+
+
+ org.assertj
+ assertj-core
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.logging.log4j
+ log4j-api
+
+
+ org.apache.logging.log4j
+ log4j-core
+
+
+ org.apache.logging.log4j
+ log4j-slf4j-impl
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+
+ org.assertj
+ assertj-core
+ 3.27.7
+ test
org.apache.flink
@@ -103,6 +214,29 @@
${flink.version}
test
tests
+
+
+ com.esotericsoftware.kryo
+ kryo
+
+
+ org.apache.commons
+ commons-lang3
+
+
+ org.lz4
+ lz4-java
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
com.github.codemonstur
@@ -116,6 +250,12 @@
${flink.version}
test
tests
+
+
+ org.xerial.snappy
+ snappy-java
+
+
org.scalatest
@@ -138,14 +278,57 @@
io.github.embeddedkafka
embedded-kafka_2.12
- 3.4.0
+ 3.9.1
test
+
+
+ log4j
+ log4j
+
+
+ org.slf4j
+ slf4j-log4j12
+
+
+ org.xerial.snappy
+ snappy-java
+
+
+ org.apache.zookeeper
+ zookeeper
+
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+ test
+
+
+
+ org.xerial.snappy
+ snappy-java
+
+
io.zonky.test
embedded-postgres
2.0.3
test
+
+
+ org.apache.commons
+ commons-compress
+
+
+
+
+ org.apache.commons
+ commons-compress
+ 1.26.0
+ test
@@ -213,7 +396,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}
@@ -250,7 +433,7 @@
org.scalatest
scalatest-maven-plugin
- 1.0
+ 2.2.0
${project.build.directory}/surefire-reports
.
@@ -259,6 +442,7 @@
test
+ test
test
diff --git a/pom.xml b/pom.xml
index c64ae25e..96616f92 100644
--- a/pom.xml
+++ b/pom.xml
@@ -20,21 +20,20 @@
dataset-registry
transformation-sdk
pipeline
- data-products
UTF-8
2.12
- 2.12.11
+ 2.12.18
1.20.0
1.20
- 3.7.1
+ 3.9.1
11
1.9.13
1.4.0
- 2.14.1
+ 2.25.3
@@ -61,6 +60,63 @@
+
+
+
+
+ org.postgresql
+ postgresql
+ 42.7.7
+
+
+
+ com.sun.mail
+ mailapi
+ 1.6.8
+
+
+
+ org.xerial.snappy
+ snappy-java
+ 1.1.10.5
+
+
+
+ com.fasterxml.woodstox
+ woodstox-core
+ 6.7.0
+
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.9.3
+
+
+ org.apache.zookeeper
+ zookeeper-jute
+ 3.9.3
+
+
+
+ ch.qos.logback
+ logback-classic
+ 1.5.18
+
+
+ ch.qos.logback
+ logback-core
+ 1.5.18
+
+
+
+ commons-io
+ commons-io
+ 2.15.1
+
+
+
+
@@ -86,4 +142,4 @@
-
+
\ No newline at end of file
diff --git a/transformation-sdk/pom.xml b/transformation-sdk/pom.xml
index 409bc554..e3fa5c64 100644
--- a/transformation-sdk/pom.xml
+++ b/transformation-sdk/pom.xml
@@ -45,8 +45,19 @@
com.fasterxml.jackson.core
jackson-annotations
+
+
+ com.fasterxml.woodstox
+ woodstox-core
+
+
+
+ com.fasterxml.woodstox
+ woodstox-core
+ 6.7.0
+
com.github.bancolombia
data-mask-core
@@ -56,7 +67,16 @@
com.fasterxml.jackson.core
jackson-databind
-
+
+ org.apache.commons
+ commons-lang3
+
+
+
+
+ org.apache.commons
+ commons-lang3
+ 3.18.0
org.scalatest
@@ -153,7 +173,7 @@
net.alchim31.maven
scala-maven-plugin
- 4.4.0
+ 4.8.1
${java.target.runtime}
${java.target.runtime}