暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Apache Flink 漫谈系列 - Flink Forward Asia 2020 盛宴之后 - PyFlink开发和调试

孙金城 2020-12-13
614

听PyFlink核心技术

开发环境依赖

PyFlink作业的开发和运行需要依赖Python 3.5/3.6/3.7 版本和Java 8或者Java 11,本游乐场所使用的环境是Java 1.8.0_211, Python 3.7.7 还有一些其他基础软件如下:

  • PyFlink 1.11

  • Java 1.8.0_211

  • Python 3.7.7

  • PIP 20.0.2

  • PyCharm Runtime version: 11.0.7

  • MocOS 10.14.6

PyCharm 配置 Python interpreter

应用PyCharm进行开发首先要配置一下项目所使用的Python环境,配置路径PyCharm -> Preferences -> Project Interpreter
如下:

点击 Add
 配置新的环境,如下:


一路”OK“,完成配置。

安装PyFlink

我们先利用PyCharm创建一些项目,名为PyFlinkPlayground
, 并为项目选择我们刚才创建的Virtualenv环境,如下:

创建之后,我们会看到External Libraries
 里面使用了PlaygroundEnv
, 但是初始化并没有PyFlink,所以我们需要进行显示的安装,如下:

我们可以手工安装PyFlink,直接在PyCharm的Terminal
下进行安装,这时候我们自动就是启动的PlaygroundEnv
环境,在安装的过程中你也可以看到site-packages
内容会不断增加,

(PlaygroundEnv) jincheng:~ jincheng.sunjc$ python --version
Python 3.7.7
(PlaygroundEnv) jincheng:~ jincheng.sunjc$ python -m pip install apache-flink==1.11.1
Collecting apache-flink==1.11.1
Using cached apache_flink-1.11.1-cp37-cp37m-macosx_10_9_x86_64.whl (206.7 MB)

...
...
Successfully installed apache-beam-2.19.0 apache-flink-1.11.1 avro-python3-1.9.1 certifi-2020.6.20 chardet-3.0.4 cloudpickle-1.2.2 crcmod-1.7 dill-0.3.1.1 docopt-0.6.2 fastavro-0.21.24 future-0.18.2 grpcio-1.30.0 hdfs-2.5.8 httplib2-0.12.0 idna-2.10 jsonpickle-1.2 mock-2.0.0 numpy-1.19.1 oauth2client-3.0.0 pandas-0.25.3 pbr-5.4.5 protobuf-3.12.4 py4j-0.10.8.1 pyarrow-0.15.1 pyasn1-0.4.8 pyasn1-modules-0.2.8 pydot-1.4.1 pymongo-3.11.0 pyparsing-2.4.7 python-dateutil-2.8.0 pytz-2020.1 requests-2.24.0 rsa-4.6 six-1.15.0 typing-3.7.4.3 typing-extensions-3.7.4.2 urllib3-1.25.10
(PlaygroundEnv) jincheng:~ jincheng.sunjc$


最终完成之后你可以在 site-packages
下面找的 pyflink
目录,如下:

有了这些信息我们就可以进行PyFlink的作业开发了。

HelloWorld 示例

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import EnvironmentSettings, StreamTableEnvironment

def hello_world():
"""
从随机Source读取数据,然后直接利用PrintSink输出。
"""
settings = EnvironmentSettings.new_instance().in_streaming_mode().use_blink_planner().build()
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(stream_execution_environment=env, environment_settings=settings)
source_ddl = """
CREATE TABLE random_source (
f_sequence INT,
f_random INT,
f_random_str STRING
) WITH (
'connector' = 'datagen',
'rows-per-second'='5',
'fields.f_sequence.kind'='sequence',
'fields.f_sequence.start'='1',
'fields.f_sequence.end'='1000',
'fields.f_random.min'='1',
'fields.f_random.max'='1000',
'fields.f_random_str.length'='10'
)
"""

sink_ddl = """
CREATE TABLE print_sink (
f_sequence INT,
f_random INT,
f_random_str STRING
) WITH (
'connector' = 'print'
)
"""

# 注册source和sink
t_env.execute_sql(source_ddl)
t_env.execute_sql(sink_ddl)

# 数据提取
tab = t_env.from_path("random_source")
# 这里我们暂时先使用 标注了 deprecated 的API, 因为新的异步提交测试有待改进...
tab.insert_into("print_sink")
# 执行作业
t_env.execute("Flink Hello World")

if __name__ == '__main__':
hello_world()


上面代码在PyCharm里面右键运行就应该打印如下结果了:

开发日志

正常来讲我们可能开发一些UDF,可能打印一些日志或者特殊情况还可能进行Python代码的调试,怎么解?

  • 首先,我们定义一个UDF,在UDF里面添加调试日志,如下:

# 定义UDF
@udf(input_types=[DataTypes.STRING()], result_type=DataTypes.STRING())
def pass_by(str):
logging.error("Some debugging infomation...")
return str


  • 然后在SQL里面使用这个UDF,如下:

# 注册 UDF
t_env.register_function('pass_by', pass_by)
# 使用UDF
tab.select("f_sequence, f_random, pass_by(f_random_str) ")


  • 完整的代码

from pyflink.datastream import StreamExecutionEnvironment
from pyflink.table import EnvironmentSettings, StreamTableEnvironment, DataTypes
from pyflink.table.udf import udf

import logging

# 定义UDF
@udf(input_types=[DataTypes.STRING()], result_type=DataTypes.STRING())
def pass_by(str):
logging.error("Some debugging infomation...")
return "pass_by_" + str

def hello_world():
"""
从随机Source读取数据,然后直接利用PrintSink输出。
"""
settings = EnvironmentSettings.new_instance().in_streaming_mode().use_blink_planner().build()
env = StreamExecutionEnvironment.get_execution_environment()
t_env = StreamTableEnvironment.create(stream_execution_environment=env, environment_settings=settings)
t_env.get_config().get_configuration().set_boolean("python.fn-execution.memory.managed", True)

source_ddl = """
CREATE TABLE random_source (
f_sequence INT,
f_random INT,
f_random_str STRING
) WITH (
'connector' = 'datagen',
'rows-per-second'='5',
'fields.f_sequence.kind'='sequence',
'fields.f_sequence.start'='1',
'fields.f_sequence.end'='1000',
'fields.f_random.min'='1',
'fields.f_random.max'='1000',
'fields.f_random_str.length'='10'
)
"""

sink_ddl = """
CREATE TABLE print_sink (
f_sequence INT,
f_random INT,
f_random_str STRING
) WITH (
'connector' = 'print'
)
"""

# 注册source和sink
t_env.execute_sql(source_ddl)
t_env.execute_sql(sink_ddl)

# 注册 UDF
t_env.register_function('pass_by', pass_by)

# 数据提取
tab = t_env.from_path("random_source")
# 这里我们暂时先使用 标注了 deprecated 的API, 因为新的异步提交测试有待改进...
tab.select("f_sequence, f_random, pass_by(f_random_str) ").insert_into("print_sink")
# 执行作业
t_env.execute("Flink Hello World")

if __name__ == '__main__':
hello_world()


那么运行之后,日志在哪里呢?就是在项目的 PlaygroundEnv -> site-packages -> pyflink -> log
 目录 ,如下:

到这里,简单的 开发环境就OK了,大家可以改改代码,直观体验一下。。。

代码调试

直观的体验之后,难免有高级用户期望单步调试自己的UDF哈,那么咱们在设置一下在PyCharm里面进行单步提示PythonUDF。

  • 设置断点端口 为我们的"HelloWord.py"设置调试端口,如下:  点击“Edit Configurations”, 添加 “Python debug server”,具体如下:  Copy 如下信息到自己的UDF:

import pydevd_pycharm
pydevd_pycharm.settrace('localhost', port=6789, stdoutToServer=True, stderrToServer=True)


这里别忘记,pip install pydevd_pycharm
, 如下:

安装pydevd_pycharm
 之后,添加debug命令到UDF,如下:

然后,点击“debug”,再执行“HelloWord", 如下:

如果一切顺利,那么程序启动之后,我们将看到断点成功,可以查看UDF的输入变量值了:)


期望你PyFlink之旅一切顺利...,知道怎会开发调试,还想知道核心原理,戳下面:


邀你入伙阿里Lemming

我们说IoT是未来10年的技术方向,新的技术方向在带来社会利好的同时,一定会潜在着新的社会问题,比如,如何让IoT中500多亿设备产生的海量(70+ZB)数据,发挥价值?

面对这样的问题,本质上需要用 云+边+端 的一体化计算引擎进行数据分析来让数据产生价值。Lemming(旅鼠)要结合云上的计算能力(Flink/MaxCompute/Holo等)构建云+边+端的计算分析能力。规划打造 MaxCompute-Things Analytics 产品,我们现在叫 Lemming(旅鼠)提到Lemming(旅鼠),大家也许更加熟悉旅鼠效应。那么,关于Lemming(旅鼠)的跳海/跳崖,在MaxCompute-Things Analytics 产品中就是上云,端和边无法承载IoT中爆发式的数据增长,让海量数据产生其应有的价值的最好办法就是上云,形成云+边+端一体的计算分析体系。以前 Lemming说:“Does not hibernate during the bitter Arctic winter!”,当Lemming跳到互联网,心系云端会说:“万物互联,弱网环境,保持数据有序上云!”。

好的,我们要营造IoT时代云+边+端的旅鼠效应(MaxCompute-Things Analytics),所以,孙金城 邀你加入互联网下半场的战斗!专注IoT数据计算分析技术 ,始终以解决社会问题为使命,不断努力...!开源+闭源,个人影响力与社会价值感并存,我们需要你的加入,如果你期望自己的明天比今天更好,那么,改变,就是迎创明天的第一步...

微信私聊(阿里招聘)

文章转载自孙金城,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论