すべてのプロダクト
Search
ドキュメントセンター

Realtime Compute for Apache Flink:Python ジョブの開発

最終更新日:Jul 18, 2026

このトピックでは、Flink Python API ジョブに関する背景情報、制限事項、開発方法、デバッグ方法、およびコネクタの使用方法について説明します。

背景情報

Flink Python ジョブはローカルで開発する必要があります。開発が完了したら、Flink 開発コンソールにジョブをデプロイして実行し、ビジネス結果を確認します。全体的な手順の詳細については、「Flink Python ジョブ」をご参照ください。

開発環境要件

  • Realtime Compute for Apache Flink の VVR 8.0.11 より前のエンジンバージョンには、Python 3.7.9 がプリインストールされています。VVR 8.0.11 以降には Python 3.9.21 がプリインストールされています。VVR 11.7 以降には、Python 3.9.21、3.10.19、および 3.11.14 がプリインストールされています。

    説明

    ローカル開発環境の Python バージョンは、対象の VVR エンジンにプリインストールされている Python バージョンと一致させることを推奨します。

    説明

    VVR 11.7 以降では、デフォルトで Python 3.9 が使用されます。プリインストール済みの Python 3.10 または 3.11 を使用する場合は、以下のようにジョブを構成してください。

    1. O&M Center > Job O&M ページで、対象のジョブ名をクリックします。

    2. Deployment Details タブの Running Parameters セクションで、右側の Edit をクリックし、Additional Configuration フィールドに次の構成を追加します(Python 3.11 を例としています)。

      python.executable: python3.11
      python.client.executable: python3.11
  • PyFlink がインストールされており、そのバージョンが対象の VVR エンジンバージョンと一致している必要があります。たとえば、デプロイページで選択したエンジンが vvr-11.7.0-jdk11-flink-1.20 の場合、次をインストールする必要があります。

    pip install ververica-flink==11.7.0
    説明

    VVR 11.5 以前では、専用の PyFlink パッケージは提供されていません。オープンソースの PyFlink の対応バージョンをインストールする必要があります。たとえば、デプロイページで選択したエンジンが vvr-8.0.9-flink-1.17 の場合、apache-flink==1.17.* をインストールする必要があります。

    pip install apache-flink==1.17.2
  • IDE がインストールされている必要があります。PyCharm または VS Code を推奨します。

  • Python ジョブはローカルで開発し、その後 Realtime Compute for Apache Flink コンソールにデプロイして実行する必要があります。

制限事項

Flink はデプロイ環境およびネットワーク環境の影響を受けるため、Python ジョブを開発する際は以下の制限事項にご注意ください。

  • オープンソース Flink V1.13 以降のみサポートされます。

  • Flink ワークスペースには、Pandas、NumPy、PyArrow などの一般的な Python ライブラリを含む Python 環境がプリインストールされています。詳細については、本トピック末尾の「プリインストール済みパッケージ一覧」をご参照ください。

  • Flink 実行環境では JDK 8 および JDK 11 のみサポートされます。Python ジョブがサードパーティの JAR パッケージに依存する場合、その JAR パッケージが互換性を持つことを確認してください。

  • VVR 4.x ではオープンソース Scala V2.11 のみサポートされます。VVR 6.x 以降ではオープンソース Scala V2.12 のみサポートされます。Python ジョブがサードパーティの JAR パッケージに依存する場合、JAR パッケージの依存関係が対応する Scala バージョンと一致することを確認してください。

ジョブ開発

Table API/SQL と DataStream API の選択

PyFlink では、Table API/SQL および DataStream API の 2 つの開発アプローチがサポートされています。Table API/SQL を推奨します。理由は以下のとおりです。

  • パフォーマンスが優れている:Table API/SQL の最適化された実行計画は JVM 内で完全に実行されます。一方、DataStream API では JVM と Python プロセス間で行単位のシリアル化および逆シリアル化が必要となり、大きなパフォーマンスオーバーヘッドが発生します。

  • 機能がより充実している:Table API/SQL はコネクタ、データ形式、ウィンドウ関数に対するより完全なサポートを提供し、SQL ジョブと同じコネクタエコシステムを共有します。

  • コミュニティでの推奨:Apache Flink コミュニティでは、PyFlink の主要な開発方向として Table API/SQL を優先しています。

SQL で表現できない複雑なカスタムロジックのシナリオでのみ、DataStream API を使用してください。

開発リファレンス

Flink ビジネスコードをローカルで開発する際は、以下のドキュメントを参照できます。開発が完了したら、コードを Flink 開発コンソールにアップロードしてジョブをデプロイします。

  • Apache Flink V1.20 のビジネスコード開発については、「Flink Python API 開発ガイド」をご参照ください。

  • Apache Flink コーディング中に発生する問題とその解決策については、「よくある質問」をご参照ください。

プロジェクト構造

Python ジョブの推奨プロジェクト構造は次のとおりです。

my-flink-python-project/
├── my_job.py                # 主ジョブファイル
├── udfs.py                  # ユーザー定義関数(任意)
├── requirements.txt         # サードパーティ Python 依存関係(任意)
└── config.properties        # 設定ファイル(任意)

依存関係管理

Python ジョブでカスタム Python 仮想環境、サードパーティ Python パッケージ、JAR ファイル、およびデータファイルを使用する方法については、「Python 依存関係の使用」をご参照ください。

ユーザー定義関数 (UDF)

以下は、機密文字列データをマスクする Python UDSF の開発例です。

from pyflink.table import DataTypes
from pyflink.table.udf import udf

@udf(result_type=DataTypes.STRING())
def mask_phone(phone: str):
    """電話番号のマスキング:先頭3桁と末尾4桁を残し、中央を****に置き換え"""
    if phone is None or len(phone) != 11:
        return phone
    return phone[:3] + '****' + phone[7:]

この UDF を 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;

UDF の登録、更新、および削除方法については、「UDF の管理」をご参照ください。

コネクタの使用

Flink でサポートされているコネクタの一覧については、「サポートされているコネクタ」をご参照ください。コネクタを使用するには、次の手順を実行します。

  1. Realtime Compute コンソール にログインします。

  2. 対象のワークスペースの Actions 列で Console をクリックします。

  3. 左側のナビゲーションウィンドウで File Management をクリックします。

  4. Upload Resource をクリックし、対象コネクタの Python パッケージを選択してアップロードします。

    独自開発のコネクタまたは Flink が提供するコネクタをアップロードできます。Flink が提供するコネクタの公式 Python パッケージをダウンロードするには、「コネクタ一覧」をご参照ください。

  5. O&M Center > Job O&M ページで、Create Deployment > Python Deployment をクリックします。Additional Dependencies フィールドで対象コネクタの Python パッケージを選択し、他のパラメーターを設定してジョブをデプロイします。

  6. デプロイされたジョブの名前をクリックします。Deployment Details タブの Running Parameters セクションで Edit をクリックします。Additional Configuration に Python コネクタパッケージの場所情報を追加します。

    ジョブが複数のコネクタ Python パッケージに依存する場合(たとえば、connector-1.jar および connector-2.jar という 2 つのパッケージ)、構成は次のようになります。

    pipeline.classpaths: 'file:///flink/usrlib/connector-1.jar;file:///flink/usrlib/connector-2.jar'
  7. ビルトインコネクタ、データ形式、およびカタログ(VVR 11.2 以降のみ)を使用するには、ジョブの Running Parameters セクションの Additional Configuration に構成を追加します。例を以下に示します。

    ## 複数のコネクタを使用
    pipeline.used-builtin-connectors: kafka;sls
    ## 転送データの複数のデータ形式
    pipeline.used-builtin-formats: avro;parquet
    ## 複数の作成済みカタログを使用
    pipeline.used-builtin-catalogs: catalogname1;catalogname2

コネクタの使用方法の詳細については、「完全なサンプルコード」をご参照ください。

ジョブデバッグ

Python UDF のコード実装でロギングを使用して、トラブルシューティング用のログ情報を出力できます。例を以下に示します。

import logging

@udf(result_type=DataTypes.BIGINT())
def add(i, j):
  logging.info("hello world")
  return i + j

ログを出力した後は、TaskManager のログファイルでログを確認できます。

ローカルデバッグ

Realtime Compute for Apache Flink はデフォルトでインターネットアクセスができないため、ローカルでテストする際にコードがオンラインデータソースに直接接続できない可能性があります。ローカルデバッグには以下のアプローチを推奨します。

  • 単体テスト:UDF に対して独立した単体テストを実行し、関数ロジックが正しいことを確認します。

  • ローカル実行:ファイルやインメモリデータなどのローカルデータソースを使用して入力をシミュレートし、ローカルでジョブを実行して処理ロジックを検証します。例を以下に示します。

    from pyflink.datastream import StreamExecutionEnvironment
    
    env = StreamExecutionEnvironment.get_execution_environment()
    # ローカルデータソースを使用してテスト
    ds = env.from_collection([('Alice', 1), ('Bob', 2), ('Alice', 3)])
    ds.key_by(lambda x: x[0]).sum(1).print()
    env.execute("local_test")
  • リモートデバッグ:オンラインデータソースに接続してデバッグするには、「コネクタを使用したジョブのローカル実行とデバッグ」をご参照ください。

ジョブデプロイメント

Python ジョブの開発が完了したら、ジョブを Realtime Compute コンソールにアップロードしてデプロイします。次の手順を実行してください。

  1. Realtime Compute コンソール にログインし、対象のワークスペースに移動します。

  2. 左側のナビゲーションウィンドウで File Management をクリックし、Python ジョブファイル(.py または .zip)をアップロードします。サードパーティの依存関係や設定ファイルがある場合は、それらもアップロードしてください。

  3. O&M Center > Job O&M ページで、Create Deployment > Python Deployment をクリックし、デプロイメント情報を入力します。

    パラメーター

    説明

    Python File Path

    アップロード済みの Python ジョブファイルを選択します。

    Entry Module

    ジョブファイルが .py ファイルの場合、このフィールドは不要です。ジョブファイルが .zip ファイルの場合、エントリモジュール名(例:my_job)を指定します。

    Additional Dependencies

    該当する場合は、コネクタ JAR ファイルまたは設定ファイルを選択します。

    Python Libraries

    該当する場合は、サードパーティ Python パッケージ(.whl または .zip)を選択します。

    Python Archives

    該当する場合は、カスタム Python 仮想環境(.zip)を選択します。

  4. Deploy をクリックします。

デプロイメントパラメーターの詳細については、「ジョブのデプロイ」をご参照ください。

完全なサンプルコード

この例では、Kafka からデータを読み取り、簡単な処理を実行し、結果を MySQL に書き込む Python ストリーミングジョブを示します。これは参考用です。

説明

この例には、チェックポイントや再起動戦略などの実行パラメーターの構成は含まれていません。これらの構成は、ジョブをデプロイした後に Deployment Details ページでカスタマイズできます。詳細については、「ジョブデプロイメント情報の構成」をご参照ください。

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():
    # 実行環境を作成
    env = StreamExecutionEnvironment.get_execution_environment()
    t_env = StreamTableEnvironment.create(env)

    # Kafka ソーステーブルを作成
    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'
        )
    """)

    # MySQL 結果テーブルを作成
    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'
        )
    """)

    # スコアが 60 以上のレコードをフィルタリングして 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()

プリインストール済みパッケージ

VVR-11

Flink ワークスペースにプリインストールされているパッケージは次のとおりです。

パッケージ

バージョン

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

Flink ワークスペースにプリインストールされているパッケージは次のとおりです。

パッケージ

バージョン

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

Flink ワークスペースにプリインストールされているパッケージは次のとおりです。

パッケージ

バージョン

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

参考資料

  • Flink Python ジョブの完全な開発ワークフローの例については、「Flink Python ジョブ」をご参照ください。

  • Flink Python ジョブでカスタム Python 仮想環境、サードパーティ Python パッケージ、JAR ファイル、およびデータファイルを使用する方法については、「Python 依存関係の使用」をご参照ください。

  • Realtime Compute for Apache Flink では、SQL ジョブおよび DataStream ジョブもサポートされています。詳細については、「ジョブ開発ガイド」および「JAR ジョブの開発」をご参照ください。