From 494d32c56531cff7b1a98cecf4cc2fe8a32a53cc Mon Sep 17 00:00:00 2001 From: Mohammadreza Khani Date: Wed, 22 Jul 2026 10:48:53 +0330 Subject: [PATCH] feat: update flink to 2.3 --- Dockerfile.flink | 63 +++++++++++++++++++++----- go.mod | 2 +- go.sum | 2 - helm/chart/Chart.yaml | 2 +- helm/chart/templates/flink/config.yaml | 8 ++++ helm/chart/values.yaml | 8 +++- internal/managed_job/run.go | 2 +- scripts/build-flink-images.sh | 58 ++++++++++++++++++++++++ start-cluster.sh | 0 9 files changed, 127 insertions(+), 18 deletions(-) create mode 100755 scripts/build-flink-images.sh mode change 100644 => 100755 start-cluster.sh diff --git a/Dockerfile.flink b/Dockerfile.flink index dd83c25..91cb6b0 100644 --- a/Dockerfile.flink +++ b/Dockerfile.flink @@ -1,33 +1,74 @@ -FROM public.ecr.aws/docker/library/flink:1.20.1-scala_2.12-java17 +ARG FLINK_VERSION=2.3.0 +ARG SCALA_VERSION=2.12 +ARG JAVA_VERSION=21 + +FROM flink:${FLINK_VERSION}-scala_${SCALA_VERSION}-java${JAVA_VERSION} # Set working directory WORKDIR /opt/flink # Set environment variables for Flink mini-cluster -ENV FLINK_HOME /opt/flink +ENV FLINK_HOME=/opt/flink ENV PATH=$FLINK_HOME/bin:$PATH # Expose necessary ports for the Flink UI (JobManager) and job manager EXPOSE 8081 6123 COPY ./start-cluster.sh /opt/flink/bin/start-cluster.sh -RUN chmod +x /opt/flink/bin/start-cluster.sh -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-connector-kafka/3.4.0-1.20/flink-connector-kafka-3.4.0-1.20.jar -P /opt/flink/lib/ +# ---- Flink-version-dependent connectors ---- + +# Kafka connector +ARG KAFKA_CONNECTOR_VERSION=5.0.0-2.2 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-connector-kafka/${KAFKA_CONNECTOR_VERSION}/flink-connector-kafka-${KAFKA_CONNECTOR_VERSION}.jar -P /opt/flink/lib/ + +# Kafka clients (version-independent) RUN wget -q https://repo1.maven.org/maven2/org/apache/kafka/kafka-clients/3.9.0/kafka-clients-3.9.0.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-avro/1.20.1/flink-avro-1.20.1.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-avro-confluent-registry/1.20.1/flink-avro-confluent-registry-1.20.1.jar -P /opt/flink/lib/ + +# Avro format +ARG AVRO_VERSION=2.2.1 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-avro/${AVRO_VERSION}/flink-avro-${AVRO_VERSION}.jar -P /opt/flink/lib/ + +# Avro Confluent Registry +ARG AVRO_REGISTRY_VERSION=2.2.1 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-avro-confluent-registry/${AVRO_REGISTRY_VERSION}/flink-avro-confluent-registry-${AVRO_REGISTRY_VERSION}.jar -P /opt/flink/lib/ + +# ---- Version-independent connectors ---- + +# ClickHouse SQL connector RUN wget -q https://repo1.maven.org/maven2/name/nkonev/flink/flink-sql-connector-clickhouse/1.17.1-8/flink-sql-connector-clickhouse-1.17.1-8.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-postgres-cdc/3.2.1/flink-sql-connector-postgres-cdc-3.2.1.jar -P /opt/flink/lib/ + +# PostgreSQL CDC +ARG PG_CDC_VERSION=3.6.0-2.2 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-postgres-cdc/${PG_CDC_VERSION}/flink-sql-connector-postgres-cdc-${PG_CDC_VERSION}.jar -P /opt/flink/lib/ + +# Legacy Avro RUN wget -q https://repo1.maven.org/maven2/org/apache/avro/avro/1.8.2/avro-1.8.2.jar -P /opt/flink/lib/ + +# misc libs RUN wget -q https://repo1.maven.org/maven2/net/objecthunter/exp4j/0.4.5/exp4j-0.4.5.jar -P /opt/flink/lib/ RUN wget -q https://jdbc.postgresql.org/download/postgresql-42.7.4.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-sql-jdbc-driver/1.20.1/flink-sql-jdbc-driver-1.20.1.jar -P /opt/flink/lib/ + +# JDBC driver +ARG JDBC_DRIVER_VERSION=2.2.1 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-sql-jdbc-driver/${JDBC_DRIVER_VERSION}/flink-sql-jdbc-driver-${JDBC_DRIVER_VERSION}.jar -P /opt/flink/lib/ + +# Legacy Scala JDBC connector (no 2.x version exists) RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-jdbc_2.12/1.10.3/flink-jdbc_2.12-1.10.3.jar -P /opt/flink/lib/ + +# JDBC connector (latest: 3.3.0-1.20, no 2.x version yet) +ARG JDBC_CONNECTOR_VERSION=3.3.0-1.20 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-connector-jdbc/${JDBC_CONNECTOR_VERSION}/flink-connector-jdbc-${JDBC_CONNECTOR_VERSION}.jar -P /opt/flink/lib/ + +# misc libs RUN wget -q https://repo1.maven.org/maven2/com/aventrix/jnanoid/jnanoid/2.0.0/jnanoid-2.0.0.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-connector-jdbc/3.2.0-1.19/flink-connector-jdbc-3.2.0-1.19.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-s3-fs-presto/1.20.1/flink-s3-fs-presto-1.20.1.jar -P /opt/flink/lib/ -RUN wget -q https://repo1.maven.org/maven2/com/github/luben/zstd-jni/1.5.7-2/zstd-jni-1.5.7-2.jar -P /opt/flink/lib/ + +# S3 filesystem +ARG S3_FS_VERSION=2.2.1 +RUN wget -q https://repo1.maven.org/maven2/org/apache/flink/flink-s3-fs-presto/${S3_FS_VERSION}/flink-s3-fs-presto-${S3_FS_VERSION}.jar -P /opt/flink/lib/ + +# Compression +RUN wget -q https://repo1.maven.org/maven2/com/github/luben/zstd-jni/1.5.7-2/zstd-jni-1.5.7-2.jar -P /opt/flink/lib/ # Command to start Flink JobManager and TaskManager in a mini-cluster setup CMD ["bin/start-cluster.sh"] \ No newline at end of file diff --git a/go.mod b/go.mod index 8d95d45..987bada 100644 --- a/go.mod +++ b/go.mod @@ -5,7 +5,7 @@ go 1.23.2 require ( github.com/danielgtaylor/huma/v2 v2.27.0 github.com/gofiber/fiber/v2 v2.52.6 - github.com/logi-camp/go-flink-client v0.2.0 + github.com/logi-camp/go-flink-client v0.2.1 github.com/samber/lo v1.47.0 go.uber.org/zap v1.27.0 k8s.io/apimachinery v0.31.3 diff --git a/go.sum b/go.sum index cf5f2d1..b0e8fb3 100644 --- a/go.sum +++ b/go.sum @@ -62,8 +62,6 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= -github.com/logi-camp/go-flink-client v0.2.0 h1:PIyfJq7FjW28bnvemReCicIuQD7JzVgJDk2xPTZUS2s= -github.com/logi-camp/go-flink-client v0.2.0/go.mod h1:A79abedX6wGQI0FoICdZI7SRoGHj15QwMwWowgsKYFI= github.com/logi-camp/go-flink-client v0.2.1 h1:STfKamFm9+2SxxfZO3ysdFsb5MViQdThB4UHbnkUOE8= github.com/logi-camp/go-flink-client v0.2.1/go.mod h1:A79abedX6wGQI0FoICdZI7SRoGHj15QwMwWowgsKYFI= github.com/mailru/easyjson v0.7.7 h1:UGYAvKxe3sBsEDzO8ZeWOSlIQfWFlxbzLZe7hwFURr0= diff --git a/helm/chart/Chart.yaml b/helm/chart/Chart.yaml index 2644d97..ae35134 100644 --- a/helm/chart/Chart.yaml +++ b/helm/chart/Chart.yaml @@ -2,5 +2,5 @@ apiVersion: v2 name: flink-kube-operator description: Helm chart for flink kube operator type: application -version: 1.2.3 +version: 1.3.0 appVersion: "0.1.1" diff --git a/helm/chart/templates/flink/config.yaml b/helm/chart/templates/flink/config.yaml index ca72088..cdec79e 100644 --- a/helm/chart/templates/flink/config.yaml +++ b/helm/chart/templates/flink/config.yaml @@ -11,7 +11,11 @@ taskmanager.data.port: 6125 taskmanager.numberOfTaskSlots: {{ .Values.flink.taskManager.numberOfTaskSlots }} parallelism.default: {{ .Values.flink.parallelism.default }} + {{- if semverCompare ">=2.0.0" .Values.flink.version }} + state.backend.type: {{ .Values.flink.state.backend }} + {{- else }} state.backend: {{ .Values.flink.state.backend }} + {{- end }} rest.port: 8081 rootLogger.level = DEBUG rootLogger.appenderRef.console.ref = ConsoleAppender @@ -25,7 +29,11 @@ {{- else if eq .Values.flink.state.checkpoint.storageType "s3" }} state.checkpoints.dir: s3://flink/checkpoints/ {{- end }} + {{- if semverCompare ">=2.0.0" .Values.flink.version }} + state.backend.rocksdb.local-dir: /opt/flink/rocksdb + {{- else }} state.backend.rocksdb.localdir: /opt/flink/rocksdb + {{- end }} high-availability.storageDir: /opt/flink/ha {{- if eq .Values.flink.state.savepoint.storageType "filesystem" }} state.savepoints.dir: file:///opt/flink/savepoints/ diff --git a/helm/chart/values.yaml b/helm/chart/values.yaml index 47308fa..5fadded 100644 --- a/helm/chart/values.yaml +++ b/helm/chart/values.yaml @@ -113,9 +113,13 @@ affinity: {} # Global values for the Flink deployment flink: + # Flink major.minor version (semver string). Used by config template to emit + # the correct configuration key names for the target Flink version. + # Override to "1.20" for backward compatibility with Flink 1.20.x. + version: "2.3" image: - repository: lcr.logicamp.tech/library/flink - tag: 1.20.1-scala_2.12-java17-minicluster + repository: pcr.pishgaman.top/library/flink-kube-operator/flink + tag: 2.3.0-java21-v1.3.0 parallelism: default: 1 # Default parallelism for Flink jobs diff --git a/internal/managed_job/run.go b/internal/managed_job/run.go index 1e543c2..5aa48f1 100644 --- a/internal/managed_job/run.go +++ b/internal/managed_job/run.go @@ -43,7 +43,7 @@ func (job *ManagedJob) Run(restoreMode bool) error { EntryClass: job.def.Spec.EntryClass, SavepointPath: savepointPath, Parallelism: job.def.Spec.Parallelism, - ProgramArg: job.def.Spec.Args, + ProgramArgsList: job.def.Spec.Args, }) if err == nil { pkg.Logger.Info("[managed-job] [run] jar successfully ran", zap.Any("run-jar-resp", runJarResp)) diff --git a/scripts/build-flink-images.sh b/scripts/build-flink-images.sh new file mode 100755 index 0000000..95ca168 --- /dev/null +++ b/scripts/build-flink-images.sh @@ -0,0 +1,58 @@ +#!/bin/bash +set -euo pipefail + +REGISTRY="${REGISTRY:-lcr.logicamp.tech/library/flink}" + +echo "=== Building Flink 2.3.0 image (default) ===" +docker build \ + --build-arg FLINK_VERSION=2.3.0 \ + --build-arg JAVA_VERSION=21 \ + --build-arg KAFKA_CONNECTOR_VERSION=5.0.0-2.2 \ + --build-arg AVRO_VERSION=2.2.1 \ + --build-arg AVRO_REGISTRY_VERSION=2.2.1 \ + --build-arg PG_CDC_VERSION=3.6.0-2.2 \ + --build-arg JDBC_DRIVER_VERSION=2.2.1 \ + --build-arg JDBC_CONNECTOR_VERSION=3.3.0-1.20 \ + --build-arg S3_FS_VERSION=2.2.1 \ + -f Dockerfile.flink \ + -t "${REGISTRY}:2.3.0-scala_2.12-java21-minicluster" \ + . + +echo "" +echo "=== Building Flink 2.2.1 image ===" +docker build \ + --build-arg FLINK_VERSION=2.2.1 \ + --build-arg JAVA_VERSION=17 \ + --build-arg KAFKA_CONNECTOR_VERSION=5.0.0-2.2 \ + --build-arg AVRO_VERSION=2.2.1 \ + --build-arg AVRO_REGISTRY_VERSION=2.2.1 \ + --build-arg PG_CDC_VERSION=3.6.0-2.2 \ + --build-arg JDBC_DRIVER_VERSION=2.2.1 \ + --build-arg JDBC_CONNECTOR_VERSION=3.3.0-1.20 \ + --build-arg S3_FS_VERSION=2.2.1 \ + -f Dockerfile.flink \ + -t "${REGISTRY}:2.2.1-scala_2.12-java17-minicluster" \ + . + +echo "" +echo "=== Building Flink 1.20.1 image (backward compat) ===" +docker build \ + --build-arg FLINK_VERSION=1.20.1 \ + --build-arg JAVA_VERSION=17 \ + --build-arg KAFKA_CONNECTOR_VERSION=3.4.0-1.20 \ + --build-arg AVRO_VERSION=1.20.1 \ + --build-arg AVRO_REGISTRY_VERSION=1.20.1 \ + --build-arg PG_CDC_VERSION=3.2.1 \ + --build-arg JDBC_DRIVER_VERSION=1.20.1 \ + --build-arg JDBC_CONNECTOR_VERSION=3.2.0-1.19 \ + --build-arg S3_FS_VERSION=1.20.1 \ + -f Dockerfile.flink \ + -t "${REGISTRY}:1.20.1-scala_2.12-java17-minicluster" \ + . + +echo "" +echo "=== Done ===" +echo "Images built:" +echo " ${REGISTRY}:2.3.0-scala_2.12-java21-minicluster" +echo " ${REGISTRY}:2.2.1-scala_2.12-java17-minicluster" +echo " ${REGISTRY}:1.20.1-scala_2.12-java17-minicluster" diff --git a/start-cluster.sh b/start-cluster.sh old mode 100644 new mode 100755