Este tópico descreve informações básicas, limites, métodos de desenvolvimento e depuração, além do uso de conectores para jobs da API Python do Flink.
Informações básicas
Desenvolva jobs Python do Flink localmente. Após concluir o desenvolvimento, implante e inicie os jobs no console de desenvolvimento do Flink para visualizar os resultados de negócio. Para mais detalhes sobre o procedimento completo, consulte Flink Python jobs.
Requisitos do ambiente de desenvolvimento
-
Versões do mecanismo Realtime Compute for Apache Flink anteriores à VVR 8.0.11 incluem o Python 3.7.9 pré-instalado. A VVR 8.0.11 e posteriores incluem o Python 3.9.21. Já a VVR 11.7 e versões posteriores trazem pré-instalados o Python 3.9.21, 3.10.19 e 3.11.14.
NotaRecomendamos que a versão do Python no seu ambiente de desenvolvimento local corresponda à versão pré-instalada no mecanismo VVR de destino.
NotaA VVR 11.7 e versões posteriores usam o Python 3.9 por padrão. Para usar o Python 3.10 ou 3.11 pré-instalado, configure o job da seguinte forma:
Na página , clique em nome do job desejado.
-
Na aba Configuration, na seção Parameters, clique em Edit no lado direito e adicione a configuração abaixo no campo Other Configuration (exemplo com Python 3.11):
python.executable: python3.11 python.client.executable: python3.11
-
Instale o PyFlink em uma versão compatível com o mecanismo VVR de destino. Por exemplo, se o mecanismo selecionado na página de implantação for
vvr-11.7.0-jdk11-flink-1.20, instale:pip install ververica-flink==11.7.0NotaA VVR 11.5 e versões anteriores não oferecem um pacote PyFlink dedicado. Nesse caso, instale a versão correspondente do PyFlink open source. Se o mecanismo escolhido na página de implantação for
vvr-8.0.9-flink-1.17, por exemplo, instaleapache-flink==1.17.*.pip install apache-flink==1.17.2 Tenha uma IDE instalada. Recomendamos PyCharm ou VS Code.
Desenvolva os jobs Python localmente e, em seguida, implante-os e execute-os no console do Realtime Compute for Apache Flink.
Limites
Devido ao impacto dos ambientes de implantação e de rede no Flink, observe os seguintes limites ao desenvolver jobs Python:
Há suporte apenas para Flink open source V1.13 ou superior.
O workspace do Flink já possui um ambiente Python pré-instalado com bibliotecas comuns como Pandas, NumPy e PyArrow. Para mais detalhes, consulte Pre-installed package list ao final deste tópico.
O ambiente de execução do Flink suporta apenas JDK 8 e JDK 11. Caso seu job Python dependa de pacotes JAR de terceiros, garanta a compatibilidade deles.
A VVR 4.x suporta exclusivamente Scala V2.11 open source. A partir da VVR 6.x, há suporte apenas para Scala V2.12 open source. Se o seu job Python depender de pacotes JAR externos, assegure-se de que as dependências correspondam à versão adequada do Scala.
Desenvolvimento de jobs
Escolha entre Table API/SQL e DataStream API
O PyFlink oferece duas abordagens de desenvolvimento: Table API/SQL e DataStream API. Recomendamos a Table API/SQL pelos motivos a seguir:
Melhor desempenho: O plano de execução otimizado da Table API/SQL roda inteiramente dentro da JVM. Já a DataStream API exige serialização e desserialização linha a linha entre a JVM e o processo Python, o que gera sobrecarga significativa de desempenho.
Recursos mais abrangentes: A Table API/SQL oferece suporte mais completo a conectores, formatos de dados e funções de janela, compartilhando o mesmo ecossistema de conectores dos jobs SQL.
Recomendação da comunidade: A comunidade Apache Flink prioriza a Table API/SQL como principal direção de desenvolvimento para o PyFlink.
Use a DataStream API apenas em cenários com lógica personalizada complexa impossível de expressar via SQL.
Referências de desenvolvimento
Consulte os documentos abaixo para desenvolver o código de negócio do Flink localmente. Depois de finalizar o desenvolvimento, carregue o código no console de desenvolvimento do Flink e implante o job.
Para desenvolver código de negócio no Apache Flink V1.20, acesse o Guia de Desenvolvimento da API Python do Flink.
Caso enfrente problemas durante a codificação no Apache Flink, consulte o FAQ para encontrar soluções.
Estrutura do projeto
A estrutura recomendada para projetos de jobs Python é a seguinte:
my-flink-python-project/
├── my_job.py # Main job file.
├── udfs.py # User-defined functions (optional)
├── requirements.txt # Third-party Python dependencies (optional)
└── config.properties # Configuration file (optional)
Gerenciamento de dependências
Para saber como utilizar ambientes virtuais Python personalizados, pacotes Python de terceiros, arquivos JAR e arquivos de dados em jobs Python, consulte Use Python dependencies.
Funções definidas pelo usuário (UDFs)
Veja abaixo um exemplo de desenvolvimento de uma UDSF em Python que mascara dados sensíveis em strings:
from pyflink.table import DataTypes
from pyflink.table.udf import udf
@udf(result_type=DataTypes.STRING())def mask_phone(phone: str):"""Phone number masking: retains the first 3 and last 4 digits, replacing the middle with ****"""if phone is None or len(phone) != 11:return phone.
return phone[:3] + '****' + phone[7:]
Para usar essa UDF em um job SQL:
CREATE TEMPORARY FUNCTION mask_phone AS 'udfs.mask_phone' LANGUAGE PYTHON;INSERT INTO sink_table
SELECT name, mask_phone(phone) AS masked_phone
FROM source_table;
Para obter informações sobre como registrar, atualizar e excluir UDFs, consulte Manage UDFs.
Uso de conectores
Para ver a lista de conectores suportados pelo Flink, consulte Supported connectors. Siga as etapas abaixo para utilizar um conector:
Faça login no console do Realtime Compute.
Clique em Console na coluna Actions do workspace desejado.
No painel de navegação à esquerda, clique em Artifacts.
-
Clique em Upload Artifact e selecione o pacote Python do conector alvo para carregar.
É possível carregar um conector desenvolvido internamente ou um fornecido pelo Flink. Para baixar os pacotes Python oficiais dos conectores do Flink, acesse a Lista de conectores.
Na página , clique em . Selecione o pacote Python do conector desejado no campo Additional Dependencies, configure os demais parâmetros e implante o job.
-
Clique em nome do job implantado. Na aba Configuration, na seção Parameters, clique em Edit. Em Other Configuration, adicione as informações de localização dos pacotes Python dos conectores.
Se o seu job depender de vários pacotes Python de conectores, por exemplo, dois pacotes chamados connector-1.jar e connector-2.jar, use a seguinte configuração:
pipeline.classpaths: 'file:///flink/usrlib/connector-1.jar;file:///flink/usrlib/connector-2.jar' -
Para utilizar conectores integrados, formatos de dados e catálogos (apenas na VVR 11.2 e posteriores), adicione a configuração na seção Running Parameters, dentro de Additional Configuration do job. Exemplo:
## Multiple connectors used. pipeline.used-builtin-connectors: kafka;sls ## Multiple data formats for data transmission. pipeline.used-builtin-formats: avro;parquet ## Multiple previously created catalogs used. pipeline.used-builtin-catalogs: catalogname1;catalogname2
Para mais detalhes sobre o uso de conectores, consulte Complete sample code.
Depuração de jobs
Na implementação de UDFs Python, use logging para gerar informações de log que auxiliem na solução de problemas. Exemplo:
import logging
@udf(result_type=DataTypes.BIGINT())
def add(i, j):
logging.info("hello world")
return i + j
Após a geração dos logs, visualize-os nos arquivos de log do TaskManager.
Depuração local
Como o Realtime Compute for Apache Flink não tem acesso à Internet por padrão, seu código pode não conseguir se conectar diretamente a fontes de dados online para testes locais. Recomendamos as seguintes abordagens para depuração local:
Testes unitários: Execute testes unitários independentes nas UDFs para garantir a correção da lógica da função.
-
Execução local: Utilize fontes de dados locais (como arquivos ou dados em memória) para simular entradas e execute o job localmente, validando a lógica de processamento. Exemplo:
from pyflink.datastream import StreamExecutionEnvironment env = StreamExecutionEnvironment.get_execution_environment()# Use a local data source for testing. ds = env.from_collection([('Alice', 1), ('Bob', 2), ('Alice', 3)]) ds.key_by(lambda x: x[0]).sum(1).print() env.execute("local_test") Depuração remota: Para depurar conectando-se a fontes de dados online, consulte Run and debug jobs with connectors locally.
Implantação de jobs
Após concluir o desenvolvimento do job Python, carregue-o no console do Realtime Compute para implantação. Siga estas etapas:
Faça login no console do Realtime Compute e acesse o workspace desejado.
No painel de navegação à esquerda, clique em Artifacts e carregue o arquivo do job Python (.py ou .zip). Se houver dependências de terceiros ou arquivos de configuração, carregue-os também.
-
Na página , clique em e preencha as informações da implantação.
Parâmetro
Descrição
Python File Path
Selecione o arquivo do job Python carregado.
Entry Module
Campo opcional para arquivos .py. Para arquivos .zip, especifique o nome do módulo de entrada, por exemplo,
my_job.Additional Dependencies
Selecione arquivos JAR de conectores ou arquivos de configuração, se aplicável.
Python Libraries
Escolha pacotes Python de terceiros (.whl ou .zip), quando necessário.
Python Archives
Selecione um ambiente virtual Python personalizado (.zip), se for o caso.
Clique em Deploy.
Para mais informações sobre os parâmetros de implantação, consulte Deploy a job.
Código de exemplo completo
Este exemplo demonstra um job de streaming Python que lê dados do Kafka, realiza um processamento simples e grava os resultados no MySQL. Serve apenas como referência.
Este exemplo não inclui a configuração de parâmetros de execução, como checkpoints e estratégias de reinício. Personalize essas configurações na página Deployment Details após a implantação do job. Para mais detalhes, consulte Configure job deployment information.
import logging
import sys
from pyflink.common import Types
from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import StreamTableEnvironment
logging.basicConfig(stream=sys.stdout, level=logging.INFO)
def kafka_to_mysql():
# Create the execution environment.
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(env)
# Create the Kafka source table.
t_env.execute_sql("""
CREATE TABLE kafka_source (
`id` INT,
`name` STRING,
`score` INT,
`event_time` TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'student_topic',
'properties.bootstrap.servers' = 'your-kafka-broker:9092',
'properties.group.id' = 'my-group',
'scan.startup.mode' = 'latest-offset',
'format' = 'json'
)
""")
# Create the MySQL sink table.
t_env.execute_sql("""
CREATE TABLE mysql_sink (
`id` INT,
`name` STRING,
`score` INT,
PRIMARY KEY (id) NOT ENFORCED
) WITH (
'connector' = 'jdbc',
'url' = 'jdbc:mysql://your-mysql-host:3306/my_database',
'table-name' = 'student',
'username' = 'your_username',
'password' = 'your_password'
)
""")
# Filter records with scores greater than 60 and write them to MySQL.
t_env.execute_sql("""
INSERT INTO mysql_sink
SELECT id, name, score
FROM kafka_source
WHERE score >= 60
""")
if __name__ == '__main__':
kafka_to_mysql()
Pacotes pré-instalados
VVR-11
Os seguintes pacotes estão pré-instalados no workspace do Flink.
|
Pacote |
Versão |
|
apache-beam |
2.48.0 |
|
avro-python3 |
1.10.2 |
|
brotlipy |
0.7.0 |
|
certifi |
2022.12.7 |
|
cffi |
1.15.1 |
|
charset-normalizer |
2.0.4 |
|
cloudpickle |
2.2.1 |
|
conda |
22.11.1 |
|
conda-content-trust |
0.1.3 |
|
conda-package-handling |
1.9.0 |
|
crcmod |
1.7 |
|
cryptography |
38.0.1 |
|
Cython |
3.0.12 |
|
dill |
0.3.1.1 |
|
dnspython |
2.7.0 |
|
docopt |
0.6.2 |
|
exceptiongroup |
1.3.0 |
|
fastavro |
1.12.1 |
|
fasteners |
0.20 |
|
find_libpython |
0.5.0 |
|
grpcio |
1.56.2 |
|
grpcio-tools |
1.56.2 |
|
hdfs |
2.7.3 |
|
httplib2 |
0.22.0 |
|
idna |
3.4 |
|
importlib_metadata |
8.7.0 |
|
iniconfig |
2.1.0 |
|
isort |
6.1.0 |
|
numpy |
1.24.4 |
|
objsize |
0.6.1 |
|
orjson |
3.9.15 |
|
packaging |
25.0 |
|
pandas |
2.3.3 |
|
pemja |
0.5.5 |
|
pip |
22.3.1 |
|
pluggy |
1.0.0 |
|
proto-plus |
1.26.1 |
|
protobuf |
4.25.8 |
|
py-spy |
0.4.0 |
|
py4j |
0.10.9.7 |
|
pyarrow |
11.0.0 |
|
pyarrow-hotfix |
0,6 |
|
pycodestyle |
2.14.0 |
|
pycosat |
0.6.4 |
|
pycparser |
2.21 |
|
pydot |
1.4.2 |
|
pymongo |
4.15.4 |
|
pyOpenSSL |
22.0.0 |
|
pyparsing |
3.2.5 |
|
PySocks |
1.7.1 |
|
pytest |
7.4.4 |
|
python-dateutil |
2.9.0 |
|
pytz |
2025.2 |
|
regex |
2025.11.3 |
|
requests |
2.32.5 |
|
ruamel.yaml |
0.18.16 |
|
ruamel.yaml.clib |
0.2.14 |
|
setuptools |
70.0.0 |
|
six |
1.16.0 |
|
tomli |
2.3.0 |
|
toolz |
0.12.0 |
|
tqdm |
4.64.1 |
|
typing_extensions |
4.15.0 |
|
tzdata |
2025.2 |
|
urllib3 |
1.26.13 |
|
wheel |
0.38.4 |
|
zipp |
3.23.0 |
|
zstandard |
0.25.0 |
|
torch |
2.5.1 |
|
torchvision |
0.20.1 |
|
transformers |
4.57.6 |
|
opencv-python-headless |
4.10.0.84 |
|
pillow |
11.3.0 |
|
ultralytics |
8.4.66 |
|
easyocr |
1.7.2 |
|
open-clip-torch |
2.32.0 |
|
rembg |
2.0.61 |
|
onnxruntime |
1.16.3 |
|
av |
14.2.0 |
|
librosa |
0.11.0 |
|
soundfile |
0.13.1 |
|
imagehash |
4.3.2 |
|
safetensors |
0.7.0 |
|
huggingface-hub |
0.36.2 |
|
tokenizers |
0.22.2 |
|
timm |
1.0.27 |
|
scikit-learn |
1.6.1 |
|
scikit-image |
0.24.0 |
|
scipy |
1.13.1 |
|
numba |
0.60.0 |
|
llvmlite |
0.43.0 |
|
matplotlib |
3.9.4 |
|
pywavelets |
1.6.0 |
|
tifffile |
2024.8.30 |
|
shapely |
2.0.7 |
|
pymatting |
1.1.15 |
|
audioread |
3.1.0 |
|
soxr |
0.5.0.post1 |
|
joblib |
1.5.3 |
|
threadpoolctl |
3.6.0 |
|
sympy |
1.13.1 |
|
mpmath |
1.3.0 |
|
python-bidi |
0.6.10 |
|
filelock |
3.19.1 |
|
pyyaml |
6.0.3 |
|
decorator |
5.3.1 |
|
msgpack |
1.1.2 |
|
ftfy |
6.3.1 |
|
wcwidth |
0.8.1 |
|
contourpy |
1.3.0 |
|
cycler |
0.12.1 |
|
fonttools |
4.60.2 |
|
importlib-resources |
5.4.0 |
|
kiwisolver |
1.4.7 |
|
lazy-loader |
0,5 |
|
attrs |
26.1.0 |
|
jsonschema |
4.25.1 |
|
jsonschema-specifications |
2025.9.1 |
|
referencing |
0.36.2 |
|
rpds-py |
0.27.1 |
|
pooch |
1.9.0 |
|
platformdirs |
4.4.0 |
VVR-8
Os seguintes pacotes estão pré-instalados no workspace do Flink.
|
Pacote |
Versão |
|
apache-beam |
2.43.0 |
|
avro-python3 |
1.9.2.1 |
|
certifi |
2025.7.9 |
|
charset-normalizer |
3.4.2 |
|
cloudpickle |
2.2.0 |
|
crcmod |
1.7 |
|
Cython |
0.29.24 |
|
dill |
0.3.1.1 |
|
docopt |
0.6.2 |
|
fastavro |
1.4.7 |
|
fasteners |
0,19 |
|
find_libpython |
0.4.1 |
|
grpcio |
1.46.3 |
|
grpcio-tools |
1.46.3 |
|
hdfs |
2.7.3 |
|
httplib2 |
0.20.4 |
|
idna |
3.10 |
|
isort |
6.0.1 |
|
numpy |
1.21.6 |
|
objsize |
0.5.2 |
|
orjson |
3.10.18 |
|
pandas |
1.3.5 |
|
pemja |
0.3.2 |
|
pip |
22.3.1 |
|
proto-plus |
1.26.1 |
|
protobuf |
3.20.3 |
|
py4j |
0.10.9.7 |
|
pyarrow |
8.0.0 |
|
pycodestyle |
2.14.0 |
|
pydot |
1.4.2 |
|
pymongo |
3.13.0 |
|
pyparsing |
3.2.3 |
|
python-dateutil |
2.9.0 |
|
pytz |
2025.2 |
|
regex |
2024.11.6 |
|
requests |
2.32.4 |
|
setuptools |
58.1.0 |
|
six |
1.17.0 |
|
typing_extensions |
4.14.1 |
|
urllib3 |
2.5.0 |
|
wheel |
0.33.4 |
|
zstandard |
0.23.0 |
VVR-6
Os seguintes pacotes estão pré-instalados no workspace do Flink.
|
Pacote |
Versão |
|
apache-beam |
2.27.0 |
|
avro-python3 |
1.9.2.1 |
|
certifi |
2024.8.30 |
|
charset-normalizer |
3.3.2 |
|
cloudpickle |
1.2.2 |
|
crcmod |
1.7 |
|
Cython |
0.29.16 |
|
dill |
0.3.1.1 |
|
docopt |
0.6.2 |
|
fastavro |
0.23.6 |
|
future |
0.18.3 |
|
grpcio |
1.29.0 |
|
hdfs |
2.7.3 |
|
httplib2 |
0.17.4 |
|
idna |
3.8 |
|
importlib-metadata |
6.7.0 |
|
isort |
5.11.5 |
|
jsonpickle |
2.0.0 |
|
mock |
2.0.0 |
|
numpy |
1.19.5 |
|
oauth2client |
4.1.3 |
|
pandas |
1.1.5 |
|
pbr |
6.1.0 |
|
pemja |
0.1.4 |
|
pip |
20.1.1 |
|
protobuf |
3.17.3 |
|
py4j |
0.10.9.3 |
|
pyarrow |
2.0.0 |
|
pyasn1 |
0.5.1 |
|
pyasn1-modules |
0.3.0 |
|
pycodestyle |
2.10.0 |
|
pydot |
1.4.2 |
|
pymongo |
3.13.0 |
|
pyparsing |
3.1.4 |
|
python-dateutil |
2.8.0 |
|
pytz |
2024.1 |
|
requests |
2.31.0 |
|
rsa |
4.9 |
|
setuptools |
47.1.0 |
|
six |
1.16.0 |
|
typing-extensions |
3.7.4.3 |
|
urllib3 |
2.0.7 |
|
wheel |
0.42.0 |
|
zipp |
3.15.0 |
Referências
Para um exemplo completo do fluxo de trabalho de desenvolvimento de jobs Python do Flink, consulte Flink Python jobs.
Para detalhes sobre o uso de ambientes virtuais Python personalizados, pacotes Python de terceiros, arquivos JAR e arquivos de dados em jobs Python do Flink, consulte Use Python dependencies.
O Realtime Compute for Apache Flink também suporta jobs SQL e DataStream. Para mais informações, consulte Job development guide e Develop JAR jobs.