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)
还没有任何评论哟~
