在java代码中写的端口号对应的是docker内部网络端口,打包成jar包后部署在webui后也相当于在docker容器里面运行了,这时候接收不到kafka发送的消息是因为容器内部通信问题,flink集群通过上图端口访问不到kafka,需要在yml配置文件把flink集群和kafka集群用一个网络连起来

version: '3.8'
services:
  jobmanager1:
    image: flink:1.18.0
    ports:
      - "8085:8081"
    command: jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  jobmanager2:
    image: flink:1.18.0
    ports:
      - "8086:8081"
    command: jobmanager
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager2
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  taskmanager1:
    image: flink:1.18.0
    depends_on:
      - jobmanager1
      - jobmanager2
    command: taskmanager
    links:
      - "jobmanager1:jobmanager"
      - "jobmanager2:jobmanager"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  taskmanager2:
    image: flink:1.18.0
    depends_on:
      - jobmanager1
      - jobmanager2
    command: taskmanager
    links:
      - "jobmanager1:jobmanager"
      - "jobmanager2:jobmanager"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  taskmanager3:
    image: flink:1.18.0
    depends_on:
      - jobmanager1
      - jobmanager2
    command: taskmanager
    links:
      - "jobmanager1:jobmanager"
      - "jobmanager2:jobmanager"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  taskmanager4:
    image: flink:1.18.0
    depends_on:
      - jobmanager1
      - jobmanager2
    command: taskmanager
    links:
      - "jobmanager1:jobmanager"
      - "jobmanager2:jobmanager"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

  taskmanager5:
    image: flink:1.18.0
    depends_on:
      - jobmanager1
      - jobmanager2
    command: taskmanager
    links:
      - "jobmanager1:jobmanager"
      - "jobmanager2:jobmanager"
    environment:
      - JOB_MANAGER_RPC_ADDRESS=jobmanager1
    networks:
      - flink-network
      - netkafka  # 添加Kafka网络

networks:
  flink-network:
    driver: bridge
  netkafka:
    external: true  # 引用外部Kafka网络

更改配置文件后就能通过本地java代码在kafka生成模拟数据,然后flink就能接受到kafka产生的数据了

更多推荐