|
| 1 | +<!-- |
| 2 | +Licensed under the Apache License, Version 2.0 (the "License"); |
| 3 | +you may not use this file except in compliance with the License. |
| 4 | +You may obtain a copy of the License at |
| 5 | +
|
| 6 | +http://www.apache.org/licenses/LICENSE-2.0 |
| 7 | +
|
| 8 | +Unless required by applicable law or agreed to in writing, software |
| 9 | +distributed under the License is distributed on an "AS IS" BASIS, |
| 10 | +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. |
| 11 | +See the License for the specific language governing permissions and |
| 12 | +limitations under the License. |
| 13 | +--> |
| 14 | + |
| 15 | +# Running Python pipelines on a local Flink cluster |
| 16 | + |
| 17 | +This guide describes a contributor workflow for validating Python Beam pipelines |
| 18 | +against a real local Flink standalone cluster. It is useful when embedded Flink |
| 19 | +is not enough, for example when validating streaming source behavior, checkpoint |
| 20 | +boundaries, or runner-visible job state in the Flink dashboard. |
| 21 | + |
| 22 | +The commands assume a Unix shell (Linux, macOS, or WSL2 on Windows) with `curl`, |
| 23 | +`tar`, and `java` on the `PATH`. |
| 24 | + |
| 25 | +* [What this setup validates](#what-this-setup-validates) |
| 26 | +* [Prerequisites](#prerequisites) |
| 27 | +* [Start a local Flink cluster](#start-a-local-flink-cluster) |
| 28 | +* [Run a Beam Python pipeline](#run-a-beam-python-pipeline) |
| 29 | +* [Troubleshooting](#troubleshooting) |
| 30 | +* [Stop the cluster](#stop-the-cluster) |
| 31 | + |
| 32 | +## What this setup validates |
| 33 | + |
| 34 | +This setup runs three separate processes: |
| 35 | + |
| 36 | +1. A Flink standalone cluster, consisting of a JobManager and a TaskManager. |
| 37 | +1. A Beam Flink Job Server, started by the Python `FlinkRunner`. |
| 38 | +1. A Python SDK harness, using `--environment_type=LOOPBACK` for local |
| 39 | + development. |
| 40 | + |
| 41 | +The Flink dashboard at `http://localhost:8081` shows the submitted Beam jobs. |
| 42 | +This is different from embedded Flink mode, where the cluster is started only |
| 43 | +for the lifetime of one job and is not useful for manual dashboard inspection. |
| 44 | + |
| 45 | +## Prerequisites |
| 46 | + |
| 47 | +Install or prepare the following: |
| 48 | + |
| 49 | +* Docker Desktop (optional), only for the alternative method of obtaining the |
| 50 | + Flink distribution. |
| 51 | +* A Unix shell: Linux, macOS, or WSL2 on Windows. |
| 52 | +* Java 11 on the `PATH`. |
| 53 | +* A Python environment with the Beam SDK dependencies installed. |
| 54 | +* A Beam source checkout for the Python code under test. |
| 55 | +* A Flink 1.20 Job Server jar built from the same Beam checkout when validating |
| 56 | + unreleased Beam changes. |
| 57 | + |
| 58 | +For a source-built Job Server jar, run this command from the Beam checkout: |
| 59 | + |
| 60 | +```sh |
| 61 | +./gradlew :runners:flink:1.20:job-server:shadowJar |
| 62 | +``` |
| 63 | + |
| 64 | +The jar is written under: |
| 65 | + |
| 66 | +```text |
| 67 | +runners/flink/1.20/job-server/build/libs/ |
| 68 | +``` |
| 69 | + |
| 70 | +## Start a local Flink cluster |
| 71 | + |
| 72 | +Use a Flink distribution whose minor version matches a Flink version supported |
| 73 | +by your Beam version. See the [Flink Version Compatibility](https://beam.apache.org/documentation/runners/flink/#flink-version-compatibility) |
| 74 | +table in the Flink Runner documentation, and confirm the exact patch version on |
| 75 | +the [Flink downloads page](https://flink.apache.org/downloads.html). This guide |
| 76 | +uses Flink 1.20. |
| 77 | + |
| 78 | +Download and unpack the binary distribution: |
| 79 | + |
| 80 | +```sh |
| 81 | +FLINK_VERSION=1.20.1 |
| 82 | +curl -fLO "https://archive.apache.org/dist/flink/flink-${FLINK_VERSION}/flink-${FLINK_VERSION}-bin-scala_2.12.tgz" |
| 83 | +tar -xzf "flink-${FLINK_VERSION}-bin-scala_2.12.tgz" -C "$HOME" |
| 84 | +export FLINK_HOME="$HOME/flink-${FLINK_VERSION}" |
| 85 | +``` |
| 86 | + |
| 87 | +Ensure these settings exist in `$FLINK_HOME/conf/config.yaml`: |
| 88 | + |
| 89 | +```yaml |
| 90 | +jobmanager.rpc.address: localhost |
| 91 | +rest.address: localhost |
| 92 | +taskmanager.numberOfTaskSlots: 2 |
| 93 | +``` |
| 94 | +
|
| 95 | +Start the cluster. The JobManager and TaskManager run as background daemons: |
| 96 | +
|
| 97 | +```sh |
| 98 | +"$FLINK_HOME/bin/start-cluster.sh" |
| 99 | +``` |
| 100 | + |
| 101 | +Verify that the JobManager and TaskManager are available: |
| 102 | + |
| 103 | +```sh |
| 104 | +curl -fsS http://localhost:8081/overview |
| 105 | +``` |
| 106 | + |
| 107 | +Expected output includes one TaskManager and two slots: |
| 108 | + |
| 109 | +```json |
| 110 | +{"taskmanagers":1,"slots-total":2,"slots-available":2,"jobs-running":0} |
| 111 | +``` |
| 112 | + |
| 113 | +You can also open the Flink dashboard in a browser: |
| 114 | + |
| 115 | +```text |
| 116 | +http://localhost:8081 |
| 117 | +``` |
| 118 | + |
| 119 | +### Alternative: extract Flink from the Docker image |
| 120 | + |
| 121 | +If a direct download is not available, copy the distribution out of the Flink |
| 122 | +Docker image with `docker cp`: |
| 123 | + |
| 124 | +```sh |
| 125 | +docker create --name flink-dist flink:1.20 |
| 126 | +docker cp flink-dist:/opt/flink "$HOME/flink-1.20" |
| 127 | +docker rm flink-dist |
| 128 | +export FLINK_HOME="$HOME/flink-1.20" |
| 129 | +``` |
| 130 | + |
| 131 | +A distribution copied out of a Docker image can contain the container hostname in |
| 132 | +`conf/config.yaml`; see [Troubleshooting](#troubleshooting). |
| 133 | + |
| 134 | +## Run a Beam Python pipeline |
| 135 | + |
| 136 | +For local Python development, use `FlinkRunner`, point it at the standalone |
| 137 | +cluster, and use `LOOPBACK` so the Python SDK harness runs in the local process. |
| 138 | + |
| 139 | +Use a source checkout on `PYTHONPATH` when validating unreleased Python changes. |
| 140 | +Set paths for your environment: |
| 141 | + |
| 142 | +```sh |
| 143 | +export BEAM_CHECKOUT="$HOME/beam" |
| 144 | +export PYTHON="$HOME/beamenv/bin/python" |
| 145 | +export FLINK_JOB_SERVER_JAR="$(find "$BEAM_CHECKOUT/runners/flink/1.20/job-server/build/libs" \ |
| 146 | + -name 'beam-runners-flink-1.20-job-server-*.jar' | head -n 1)" |
| 147 | +``` |
| 148 | + |
| 149 | +Run a small pipeline: |
| 150 | + |
| 151 | +```sh |
| 152 | +printf 'to be or not to be\nbeam runs on flink\n' > /tmp/beam-flink-input.txt |
| 153 | + |
| 154 | +PYTHONPATH="$BEAM_CHECKOUT/sdks/python" "$PYTHON" -m apache_beam.examples.wordcount \ |
| 155 | + --runner=FlinkRunner \ |
| 156 | + --flink_master=localhost:8081 \ |
| 157 | + --flink_version=1.20 \ |
| 158 | + --flink_job_server_jar="$FLINK_JOB_SERVER_JAR" \ |
| 159 | + --environment_type=LOOPBACK \ |
| 160 | + --input=/tmp/beam-flink-input.txt \ |
| 161 | + --output=/tmp/beam-flink-counts |
| 162 | +``` |
| 163 | + |
| 164 | +For released Beam, omit `--flink_job_server_jar` and the `PYTHONPATH` prefix; the |
| 165 | +`FlinkRunner` downloads a Job Server matching `--flink_version` automatically. The |
| 166 | +source checkout and built jar are only needed to test unreleased changes. |
| 167 | + |
| 168 | +Check the dashboard or REST API after the run: |
| 169 | + |
| 170 | +```sh |
| 171 | +curl -fsS http://localhost:8081/jobs/overview |
| 172 | +``` |
| 173 | + |
| 174 | +The job should be `FINISHED`. |
| 175 | + |
| 176 | +## Troubleshooting |
| 177 | + |
| 178 | +If the TaskManager does not register, check `$FLINK_HOME/conf/config.yaml`. |
| 179 | +When a distribution is copied out of a Docker image, the file might contain the |
| 180 | +container hostname. Replace it with: |
| 181 | + |
| 182 | +```yaml |
| 183 | +jobmanager.rpc.address: localhost |
| 184 | +``` |
| 185 | +
|
| 186 | +If a Python job fails on native Windows with an invalid path containing `:`, |
| 187 | +run the Python driver and Job Server from WSL2. Some staged artifact names used |
| 188 | +by the portable runner are valid on Linux but invalid as native Windows file |
| 189 | +names. |
| 190 | + |
| 191 | +On WSL2, keep at least one shell open in the distribution while the cluster runs. |
| 192 | +Closing the last shell can stop the distribution and its background daemons. |
| 193 | + |
| 194 | +If the job starts but the Python transforms do not execute, check the |
| 195 | +environment type. `LOOPBACK` is intended for local development. For a remote |
| 196 | +or multi-machine Flink cluster, use a containerized environment instead. |
| 197 | + |
| 198 | +## Stop the cluster |
| 199 | + |
| 200 | +Stop the local cluster when you finish collecting results: |
| 201 | + |
| 202 | +```sh |
| 203 | +"$FLINK_HOME/bin/stop-cluster.sh" |
| 204 | +``` |
0 commit comments