Я могу запустить такое задание, но ничего не вижу в пользовательском интерфейсе FLink. Но из журналов докера я знаю, что задание выполнено успешно.
Я могу получить доступ только к localhost:8081, потому что менеджер заданий и диспетчер задач запущены. Но завершенное задание Python не отображается в разделе завершенных заданий пользовательского интерфейса Flink.
Вот простой код Python
Код: Выделить всё
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.common.typeinfo import Types
from pyflink.datastream.functions import MapFunction
class Splitter(MapFunction):
def map(self, value):
return value.split()
def word_count():
env = StreamExecutionEnvironment.get_execution_environment()
# Create a DataStream with sample sentences
text = env.from_collection([
'hello world',
'hello flink',
'flink is fun'
])
# Apply transformations to count words
word_stream = text.map(Splitter(), output_type=Types.STRING()) \
.flat_map(lambda words: [(word, 1) for word in words], output_type=Types.TUPLE([Types.STRING(), Types.INT()])) \
.key_by(lambda x: x[0]) \
.sum(1)
word_stream.print()
env.execute('word_count_job')
if __name__ == '__main__':
word_count()
Код: Выделить всё
# Base image for Flink
FROM apache/flink:1.18.0-scala_2.12
# Install Python 3.10 and pip
RUN apt-get update && apt-get install -y python3.10 python3-pip
# Set Python 3.10 as the default for `python` and `pip`
RUN update-alternatives --install /usr/bin/python python /usr/bin/python3.10 1 \
&& update-alternatives --install /usr/bin/pip pip /usr/bin/pip3 1
# Install PyFlink
RUN pip install apache-flink==1.18.0
# Set working directory
WORKDIR /opt/flink/python_job
# Copy the Python Flink job script into the container
COPY word_count.py /opt/flink/python_job/
# Set the entry point to submit the job to the Flink cluster using the Flink CLI
# ENTRYPOINT [ "flink", "run", "-m", "jobmanager:8081", "/opt/flink/python_job/word_count.py" ]
ENTRYPOINT [ "python", "/opt/flink/python_job/word_count.py" ]
Код: Выделить всё
version: '3.8'
services:
jobmanager:
build:
context: .
dockerfile: Dockerfile
image: flink:1.18
container_name: jobmanager
hostname: jobmanager
networks:
- flink-network
ports:
- "8081:8081" # Flink UI
# - "6123:6123" # RPC Port
expose:
- "6123"
environment:
- JOB_MANAGER_RPC_ADDRESS=jobmanager
command: jobmanager
taskmanager:
build:
context: .
dockerfile: Dockerfile
image: flink:1.18
container_name: taskmanager
hostname: taskmanager
# expose:
# - "6121"
# - "6122"
networks:
- flink-network
depends_on:
- jobmanager
links:
- jobmanager:jobmanager
environment:
- JOB_MANAGER_RPC_ADDRESS=jobmanager
- TASK_MANAGER_NUMBER_OF_TASK_SLOTS=5
command: taskmanager
pyflink-job:
build:
context: . # Path to your Dockerfile
dockerfile: Dockerfile
container_name: pyflink-job
hostname: pyflink-job
networks:
- flink-network
depends_on:
- jobmanager
- taskmanager
environment:
- JOB_MANAGER_RPC_ADDRESS=jobmanager
# entrypoint: flink run -m jobmanager:8081 /opt/flink/python_job/word_count.py
networks:
flink-network:
driver: bridge
Сообщите мне, какие изменения я могу здесь внести?
Подробнее здесь: https://stackoverflow.com/questions/790 ... ith-docker