Advertisement

Pyflink框架中的map函数应用

阅读量:

导入json模块
导入正则表达式库
导入日志记录器模块
导入系统库
从collections模块导入Counter类

import DataStream and StreamExecutionEnvironment from pyflink.datastream;
import RuntimeContext、FlatMapFunction and MapFunction from pyflink.datastream.functions;
import Types from pyflink.common.typeinfo;

获取当前流执行环境。
创建数据流实例:
使用DataStream类,
调用socketTextStream方法,
创建一个基于IP地址192.168.137.200和端口号8899的文本数据流实例。
注释说明:
显示数据流实例的创建过程,
并显示其连接到指定IP地址和服务端口的过程。
定义一个名为get_key的方法,
返回固定值'999'。
定义LogEvent类,
初始化其属性为None。

init方法接收参数world并将其赋值给实例变量self.world
#计数设置为count
get_dict方法返回包含两个键值对的字典
"world": 将自变

全部评论 (0)

还没有任何评论哟~