Skip to content

Commit 42448da

Browse files
committed
Refactor Flink Load Tests to use native Dataproc components and Auto Zone
- Configure Flink properties (off-heap memory, slots, checkpoints) natively via Dataproc --properties, resolving Flink 2.0 heartbeat timeouts and YARN allocation errors. - Delete obsolete flink.sh and docker.sh init actions, leaving only beam.sh to pull images. - Enable Dataproc Auto Zone across Flink load test workflows, removing GCLOUD_ZONE pinning. - Clean up obsolete FLINK_DOWNLOAD_URL and HADOOP_DOWNLOAD_URL parameters.
1 parent 397d9c6 commit 42448da

15 files changed

Lines changed: 54 additions & 396 deletions

.github/workflows/beam_LoadTests_Go_CoGBK_Flink_batch.yml

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,12 +47,10 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-go-cogbk-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
55-
FLINK_TASKMANAGER_SLOTS: 5
53+
FLINK_TASKMANAGER_SLOTS: 2
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
5856
JOB_SERVER_IMAGE: gcr.io/apache-beam-testing/beam_portability/beam_flink_job_server:latest-flink2.2

.github/workflows/beam_LoadTests_Go_Combine_Flink_Batch.yml

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,12 +47,10 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-go-combine-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
55-
FLINK_TASKMANAGER_SLOTS: 5
53+
FLINK_TASKMANAGER_SLOTS: 2
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
5856
JOB_SERVER_IMAGE: gcr.io/apache-beam-testing/beam_portability/beam_flink_job_server:latest-flink2.2

.github/workflows/beam_LoadTests_Go_GBK_Flink_Batch.yml

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,12 +47,10 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-go-gbk-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
55-
FLINK_TASKMANAGER_SLOTS: 5
53+
FLINK_TASKMANAGER_SLOTS: 2
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
5856
JOB_SERVER_IMAGE: gcr.io/apache-beam-testing/beam_portability/beam_flink_job_server:latest-flink2.2

.github/workflows/beam_LoadTests_Go_ParDo_Flink_Batch.yml

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-go-pardo-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest

.github/workflows/beam_LoadTests_Go_SideInput_Flink_Batch.yml

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,12 +47,10 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-go-sideinput-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
55-
FLINK_TASKMANAGER_SLOTS: 5
53+
FLINK_TASKMANAGER_SLOTS: 2
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
5856
JOB_SERVER_IMAGE: gcr.io/apache-beam-testing/beam_portability/beam_flink_job_server:latest-flink2.2

.github/workflows/beam_LoadTests_Python_CoGBK_Flink_Batch.yml

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-py-cogbk-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest

.github/workflows/beam_LoadTests_Python_Combine_Flink_Batch.yml

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-py-cmb-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
@@ -99,7 +97,7 @@ jobs:
9997
cd ${{ github.workspace }}/.test-infra/dataproc; ./flink_cluster.sh create
10098
- name: get current time
10199
run: echo "NOW_UTC=$(date '+%m%d%H%M%S' --utc)" >> $GITHUB_ENV
102-
# The env variables are created and populated in the test-arguments-action as "<github.job>_test_arguments_<argument_file_paths_index>"
100+
# The env variables are created and populated in the test-arguments-action as "<github.job>_test_arguments_<argument_file_paths_index>"
103101
- name: run Load test 2GB 10 byte records
104102
env:
105103
CLOUDSDK_CONFIG: ${{ env.KUBELET_GCLOUD_CONFIG_PATH}}
@@ -137,4 +135,4 @@ jobs:
137135
- name: Teardown Flink
138136
if: always()
139137
run: |
140-
${{ github.workspace }}/.test-infra/dataproc/flink_cluster.sh delete
138+
${{ github.workspace }}/.test-infra/dataproc/flink_cluster.sh delete

.github/workflows/beam_LoadTests_Python_Combine_Flink_Streaming.yml

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-py-cmb-flink-streaming-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest

.github/workflows/beam_LoadTests_Python_GBK_Flink_Batch.yml

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-py-gbk-flk-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
@@ -151,5 +149,5 @@ jobs:
151149
if: always()
152150
run: |
153151
${{ github.workspace }}/.test-infra/dataproc/flink_cluster.sh delete
154-
155-
# TODO(https://github.com/apache/beam/issues/20146) Re-enable auto builds after these tests pass.
152+
153+
# TODO(https://github.com/apache/beam/issues/20146) Re-enable auto builds after these tests pass.

.github/workflows/beam_LoadTests_Python_ParDo_Flink_Batch.yml

Lines changed: 2 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -47,11 +47,9 @@ env:
4747
GRADLE_ENTERPRISE_CACHE_PASSWORD: ${{ secrets.GE_CACHE_PASSWORD }}
4848
INFLUXDB_USER: ${{ secrets.INFLUXDB_USER }}
4949
INFLUXDB_USER_PASSWORD: ${{ secrets.INFLUXDB_USER_PASSWORD }}
50-
GCLOUD_ZONE: us-central1-a
50+
GCLOUD_REGION: us-central1
5151
CLUSTER_NAME: beam-loadtests-py-pardo-flink-batch-${{ github.run_id }}
5252
GCS_BUCKET: gs://beam-flink-cluster
53-
FLINK_DOWNLOAD_URL: https://archive.apache.org/dist/flink/flink-2.2.1/flink-2.2.1-bin-scala_2.12.tgz
54-
HADOOP_DOWNLOAD_URL: https://repo.maven.apache.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
5553
FLINK_TASKMANAGER_SLOTS: 1
5654
DETACHED_MODE: true
5755
HARNESS_IMAGES_TO_PULL: gcr.io/apache-beam-testing/beam-sdk/beam_go_sdk:latest
@@ -128,4 +126,4 @@ jobs:
128126
-PloadTest.mainClass=apache_beam.testing.load_tests.pardo_test \
129127
-Prunner=PortableRunner \
130128
-PpythonVersion=3.10 \
131-
'-PloadTest.args=${{ env.beam_LoadTests_Python_ParDo_Flink_Batch_test_arguments_3 }} --job_name=load-tests-python-flink-batch-pardo-4-${{ steps.datetime.outputs.datetime }}'
129+
'-PloadTest.args=${{ env.beam_LoadTests_Python_ParDo_Flink_Batch_test_arguments_3 }} --job_name=load-tests-python-flink-batch-pardo-4-${{ steps.datetime.outputs.datetime }}'

0 commit comments

Comments
 (0)