conda activate pyflink-py38env cd -/pyflink-py38env/lib/python3.8/site-packages/pyflink/ python examples/table/word_count.py -----------<log>---------------- (pyflink-py38env) [root@node80 pyflink]# python examples/table/word_count.py Using Any for unsupported type: typing.Sequence[~T] No module named google.cloud.bigquery_storage_v1. As a result, the ReadFromBigQuery transform *CANNOT* be used with `method=DIRECT_READ`. Executing word_count example with default input data set. Use --input to specify file input. Printing result to stdout. Use --output to specify output path. +I[To, 1] +I[be,, 1] +I[or, 1] +I[not, 1] +I[to, 1] +I[be,--that, 1] +I[is, 1] +I[the, 1] +I[question:--, 1] +I[Whether, 1] ...
在使用pyflink进行程序开发时,用户在开发UDF算子时常常会引入第三方依赖库,此时对依赖包引入和运行的管理就会成为刚需;当在本地执行pyflink时,用户可以将第三方python库下载安装到本地后在进行执行,但当pyflink程序需要提交到远程执行时,此方法就行不通。于是pyflink官方提供了对依赖包的管理方法,其中在DataStream API 和 Table API实现方式是不同的。
当通过flink run 提交pyflink程序时,flink将执行python命令,解析编译pyflink程序,生成一个jar程序,此过程也被称为JobGraph对象生成,其本质是python vm 与java vm 通过RPC的方式进行通信,当然,不同的提交模式,其实现逻辑会不同;目前,提交模式分为以下几种:
在Flink On Yarn高可用架构下,目前遇到提交多个任务,在集群中只有一个任务正常执行,其他任务提交成功但未执行的问题,目前解决方案是:在flink-conf.yaml配置中,将cluster.id配置内容注释掉,此配置主要用于zookeeper管理多个flink集群,至于高可用部署场景该配置是否为必配置项,暂未深入分析;参考