Hello, I am having trouble when deploying Apache Flink 1.16.1 on 2 Google Cloud instances with Docker Swarm. The JobManager is deployed on the manager node and the TaskManager is deployed in the worker node. The TaskManager seems to have trouble to communicate with ResourceManager on JobManager via akka.tcp. Here is the log on the TaskManager
[email protected] | 2023-03-29 13:50:34,061 INFO org.apache.flink.runtime.taskexecutor.DefaultJobLeaderService [] - Start job leader service. [email protected] | 2023-03-29 13:50:34,066 INFO org.apache.flink.runtime.filecache.FileCache [] - User file cache uses directory /tmp/flink-dist-cache-705fdca5-c285-4875-9b68-556ccd1b56c3 [email protected] | 2023-03-29 13:50:34,073 INFO org.apache.flink.runtime.taskexecutor.TaskExecutor [] - Connecting to ResourceManager akka.tcp://flink@jobmanager:6123/user/rpc/resourcemanager_*(00000000000000000000000000000000). [email protected] | 2023-03-29 13:50:34,420 INFO org.apache.flink.runtime.taskexecutor.TaskExecutor [] - Resolved ResourceManager address, beginning registration [email protected] | 2023-03-29 13:55:34,086 ERROR org.apache.flink.runtime.taskexecutor.TaskExecutor [] - Fatal error occurred in TaskExecutor akka.tcp://[email protected]:6127/user/rpc/taskmanager_0. [email protected] | org.apache.flink.runtime.taskexecutor.exceptions.RegistrationTimeoutException: Could not register at the ResourceManager within the specified maximum registration duration PT5M. This indicates a problem with this instance. Terminating now. About firewall rules, these instances are deployed in the same VPC, same subnet, so I suppose it can communicate without any troubles. In fact, I can ping or curl on both instances. Before deploying to Google Cloud, I have successfully deployed the same setup on my private Microstack cloud without any problem. Here is my docker-compose file that I used docker stack deploy: version: '3.8' services: jobmanager: image: halo93/fixed-ports-flink-docker:1.16.1-scala_2.12-java11-custom deploy: replicas: 1 placement: constraints: [node.hostname == managernode.gcp] ports: - "8081:8081" - "6123:6123" - "6124:6124" - "6125:6125" command: jobmanager environment: - FLINK_PROPERTIES=${FLINK_PROPERTIES} networks: - flink-network taskmanager: image: halo93/fixed-ports-flink-docker:1.16.1-scala_2.12-java11-custom deploy: replicas: 1 placement: constraints: [node.hostname == workernode.gcp] depends_on: - jobmanager ports: - "6121:6121" - "6122:6122" - "6126:6126" - "6127:6127" - "6128:6128" - "5005:5005/udp" command: - taskmanager environment: - FLINK_PROPERTIES=${FLINK_PROPERTIES} networks: - flink-network networks: flink-network: driver: overlay attachable: true FLINK_PROPERTIES FLINK_PROPERTIES=$'\njobmanager.rpc.address: jobmanager\nparallelism.default: 2\n' I am using a customized flink docker image to fix taskmanager.data.port and taskmanager.rpc.port to 6126 and 6127. I have tried to change jobmanager.rpc.address with the private IP and zonal DNS, ResourceManager can register the taskmanager. However, by doing so, flink-metrics is unable to work. I expect that I can successfully deploy flink cluster on Google Cloud instances with Docker Swarm, same with what I did on Microstack.
