大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载PyFlink 的 Table 生态系统中DataType是描述数据逻辑类型的核心抽象它贯穿于建表、声明 Python 用户自定义函数UDF输入输出类型、以及向量化 UDF 的整个流程。本文基于当前仓库中的官方文档 python_types.md 展开并结合 types.py 源码逐类剖析每种数据类型的创建方法、参数约束与底层转换机制帮助你准确声明类型、规避精度与可空性陷阱写出可正确运行的 PyFlink 作业。DataType逻辑类型与物理表示的解耦在 Table 生态系统中数据类型Data Type用于描述一个值的逻辑类型。它最常见的用途是声明 Python 用户自定义函数的输入/输出类型也用于定义表结构。PyFlink 中用户直接操作的是pyflink.table.types.DataType实例。官方文档明确强调DataType实例声明的是数据的逻辑类型这并不蕴含数据传输或存储时的具体物理表示形式。逻辑类型独立于物理表示且与 SQL 标准中的数据类型术语非常接近物理表示提示physical hints则只有在 Table 生态系统的边界例如与 DataStream API 桥接、与外部存储交互才是必需的。从源码结构看DataType 是所有类型的基类它承载了两个职责声明逻辑类型以及为优化器提供物理表示提示。基类定义了一系列通用能力可空性构造参数nullable默认True控制该类型是否允许Nonenot_null()返回一个可空性为False的副本nullable()则返回可空性为True的副本桥接提示bridged_to(conversion_cls)用于在进入或离开 Table 生态系统时提示数据应使用给定类来表示类型转换to_sql_type(obj)将 Python 对象转换为内部 SQL 对象from_sql_type(obj)反向转换而need_conversion()用于判断是否需要转换对ARRAY/MULTISET/MAP/ROW等复合类型可避免不必要的转换开销。类型体系的继承关系也从源码中得到印证AtomicType原子类型即非数组、非行、非映射的一切类型派生出NullType、NumericType、DecimalType、TimeType等而 ArrayType、MapType、MultisetType、RowType 则直接继承自DataType。创建 DataTypeDataTypes 工厂方法所有预定义的数据类型都位于pyflink.table.types包中并可通过 pyflink.table.types.DataTypes 中定义的静态工厂方法实例化。官方文档指出完整的数据类型列表可参见 docs/content/docs/dev/table/types.md 的List of Data Types章节。一个典型用法示例from pyflink.table.types import DataTypes # 原子类型 DataTypes.BOOLEAN() DataTypes.INT() DataTypes.DOUBLE() # 带长度/精度的类型 DataTypes.VARCHAR(100) DataTypes.DECIMAL(38, 18) DataTypes.TIMESTAMP(3)原子类型与数值类型从 DataTypes 的工厂方法实现可以看到各类的取值范围与约束工厂方法底层类型说明DataTypes.NULL()NullType表示无类型的None值可被转换为任意可空类型当前文档标注该类型尚未完全支持DataTypes.BOOLEAN()BooleanType布尔值SQL 标准的三值逻辑 TRUE / FALSE / UNKNOWNDataTypes.TINYINT()TinyIntType1 字节有符号整数范围 -128 ~ 127DataTypes.SMALLINT()SmallIntType2 字节有符号整数范围 -32,768 ~ 32,767DataTypes.INT()IntType4 字节有符号整数范围 -2,147,483,648 ~ 2,147,483,647DataTypes.BIGINT()BigIntType8 字节有符号整数DataTypes.FLOAT()FloatType4 字节单精度浮点数DataTypes.DOUBLE()DoubleType8 字节双精度浮点数DataTypes.DECIMAL(p, s)DecimalType固定精度与小数位的十进制数precision范围 1~38scale范围 0~precision当前实现要求 precision38、scale18字符串与二进制类型DataTypes.CHAR(length)定长字符串长度范围 1 ~ 2147483647DataTypes.VARCHAR(length)变长字符串length为最大长度范围同上当前实现中长度必须为 0x7fffffff2147483647DataTypes.STRING()DataTypes.VARCHAR(2147483647)的快捷方式DataTypes.BINARY(length)定长二进制串字节序列DataTypes.VARBINARY(length)变长二进制串DataTypes.BYTES()DataTypes.VARBINARY(2147483647)的快捷方式。时间类型DataTypes.DATE()日期年-月-日取值从0000-01-01到9999-12-31相比 SQL 标准起始年份为 0000DataTypes.TIME(precision0)无时区的时间时:分:秒[.小数]precision小数秒位数范围 0~9当前要求必须为 0DataTypes.TIMESTAMP(precision6)无时区的时间戳precision默认 6范围 0~9当前要求必须为 3。该类型不存储或表示时区描述的是挂钟上的本地时间无法单独表示时间线上的某一瞬间DataTypes.TIMESTAMP_WITH_LOCAL_TIME_ZONE(precision6)/DataTypes.TIMESTAMP_LTZ(precision6)带本地时区的时间戳内部以 long 存储全部日期时间字段纳秒精度以及相对 UTC/Greenwich 的偏移当前仅支持精度 3。区间类型INTERVALDataTypes.INTERVAL(upper_resolution, lower_resolution)用于声明时间区间分为两类Day-Time 区间由DAY()/HOUR()/MINUTE()/SECOND()组合出分辨率取值从-999999 23:59:59.999999999到999999 23:59:59.999999999例如DataTypes.INTERVAL(DataTypes.DAY(2), DataTypes.SECOND(9))Year-Month 区间由YEAR()/MONTH()组合出分辨率取值从-9999-11到9999-11例如DataTypes.INTERVAL(DataTypes.YEAR(4), DataTypes.MONTH())。YEAR(precision2)的年份位数范围 1~4DAY(precision2)的天数位数范围 1~6SECOND(precision6)的小数秒位数范围 0~9。当前实现中upper_resolution对 Year-Month 必须为MONTH、对 Day-Time 必须为SECOND且lower_resolution必须为None。复合类型DataTypes.ARRAY(element_type)元素类型相同的数组。相比 SQL 标准数组的最大基数固定为 2147483647且任意有效类型都可作为元素类型DataTypes.MAP(key_type, value_type)键到值的关联映射键不可重复MapType的键不允许为 nullDataTypes.MULTISET(element_type)多重集bag允许同一元素出现多次每个唯一值映射到一个多重度DataTypes.ROW([DataTypes.FIELD(f1, DataTypes.INT()), ...])字段序列每个字段由DataTypes.FIELD(name, data_type, descriptionNone)定义字段名 字段类型 可选描述。表的行类型即为最具体的行类型行中的每一列与行类型中序号相同的字段对应。此外源码中还提供了DataTypes.LIST_VIEW(element_type)与DataTypes.MAP_VIEW(key_type, value_type)它们只能用于聚合函数Aggregate Function的累加器类型声明不可用于常规数据。DataType 与 Python 类型、Pandas 类型的映射关系数据类型可用于声明 Python UDF 的输入/输出类型输入数据会被转换为与所声明数据类型相对应的 Python 对象UDF 执行结果的类型也必须与所声明的数据类型匹配。对于向量化 Python UDF输入类型和输出类型均为pandas.Seriespandas.Series中的元素类型与指定的数据类型对应。官方文档给出的完整映射表如下Data TypePython TypePandas TypeBOOLEANboolnumpy.bool_TINYINTintnumpy.int8SMALLINTintnumpy.int16INTintnumpy.int32BIGINTintnumpy.int64FLOATfloatnumpy.float32DOUBLEfloatnumpy.float64VARCHARstrstrVARBINARYbytesbytesDECIMALdecimal.Decimaldecimal.DecimalDATEdatetime.datedatetime.dateTIMEdatetime.timedatetime.timeTimestampTypedatetime.datetimedatetime.datetimeLocalZonedTimestampTypedatetime.datetimedatetime.datetimeINTERVAL YEAR TO MONTHint暂不支持INTERVAL DAY TO SECONDdatetime.timedelta暂不支持ARRAYlistnumpy.ndarrayMULTISETlist暂不支持MAPdict暂不支持ROWRowdict映射中的关键细节整数宽度的折叠TINYINT/SMALLINT/INT/BIGINT在 Python 侧统一映射为intPython 3 的int无固定宽度但向量化 UDF 的 Pandas 列则会按numpy.int8/int16/int32/int64区分宽度写入超过范围的值可能溢出或报错时间对象的精度差异TIMESTAMP无时区与LocalZonedTimestampType本地时区在 Python 侧都映射为datetime.datetime但由于底层存储机制不同后者以 long 存储 UTC 偏移在跨时区计算时行为不同区间类型的特别之处INTERVAL YEAR TO MONTH映射为int月数INTERVAL DAY TO SECOND映射为datetime.timedelta向量化场景下两者均暂不支持ROW的 Python 表示普通 UDF 中行映射为Row向量化 UDF 中映射为dict。Row定义于 flink-python/pyflink/common/types.py可像属性row.key或字典键row[key]一样访问字段支持as_dict(recursiveFalse)递归转为字典也可先通过Person Row(name, age)创建行结构再实例化。实战在 Python UDF 中声明与使用数据类型结合上述类型与映射规则下面给出一个完整可运行的示例演示标量 UDF 与向量化 UDF 中的类型声明from pyflink.common import Row from pyflink.table import EnvironmentSettings, TableEnvironment from pyflink.table.expressions import col from pyflink.table.udf import udf, udtf, udaf from pyflink.table.types import DataTypes import pandas as pd env_settings EnvironmentSettings.in_streaming_mode() t_env TableEnvironment.create(env_settings) # 1) 标量 UDF显式声明输入/输出类型 udf(result_typeDataTypes.STRING(), input_types[DataTypes.INT(), DataTypes.INT()]) def add(i, j): # 输入已按 INT 转换为 Python int输出 str 与 STRING 匹配 return str(i j) # 2) 向量化 UDF输入输出为 pandas.Series元素类型为 numpy.int64 / str udf(result_typeDataTypes.BIGINT(), input_types[DataTypes.BIGINT()], func_typepandas) def double(x: pd.Series) - pd.Series: return x * 2 # 3) ROW 复合类型输入为 Row可按键名访问 udf(result_typeDataTypes.ROW([DataTypes.FIELD(name, DataTypes.STRING()), DataTypes.FIELD(age, DataTypes.INT())]), input_types[DataTypes.ROW([DataTypes.FIELD(name, DataTypes.STRING()), DataTypes.FIELD(age, DataTypes.INT())])]) def grow_up(r): return Row(r.name, r.age 1) # 建表并注册 UDF 后即可在 SQL / Table API 中使用使用要点result_type与input_types中的DataType声明必须与实际 Python 值匹配例如VARCHAR对应str、ARRAY对应list、MAP对应dict、DECIMAL对应decimal.Decimal声明TIMESTAMP(3)时Python 侧传入/返回datetime.datetime对象即可框架负责与内部 SQL 对象互转向量化 UDF 需要额外传入func_typepandas且数据类型须落在上述映射表中 Pandas 类型受支持的列内如ARRAY→numpy.ndarray区间类型与MULTISET/MAP等暂不支持列则不能用于向量化声明。底层转换机制与源码印证DataType基类中to_sql_type/from_sql_type/need_conversion三个方法构成了 Python 对象与内部 SQL 对象互转的契约见 types.py原子类型通常无需转换need_conversion()返回False直接透传对象复合类型则会委托给元素类型例如 ArrayType 的to_sql_type会对list中每个元素递归调用element_type.to_sql_typefrom_sql_type则反向执行MapType 对dict的键和值分别转换。只有当元素/键值类型需要转换时need_conversion()才返回True从而在多数场景下避免无谓的拷贝开销。这一设计解释了映射表的存在意义映射表规定了Python 对象长什么样而DataType子类负责在 UDF 边界执行转换保证用户拿到的输入、以及框架校验的输出都符合声明的逻辑类型。若 UDF 返回值类型与声明的DataType不符运行时会报类型不匹配错误。小结DataType是 PyFlink Table API 中描述值逻辑类型的核心抽象与物理存储表示解耦所有预定义类型通过pyflink.table.types.DataTypes工厂方法创建完整列表见 docs/content/docs/dev/table/types.md映射表明确了每种DataType对应的 Python 类型与 Pandas 类型是编写标量 UDF 与向量化 UDF 时声明input_types/result_type的直接依据创建类型时需注意当前实现中的约束VARCHAR长度必须为 2147483647可用STRING()快捷方式、DECIMAL需 precision38/scale18、TIMESTAMP系列仅支持精度 3、MapType键不可为 null、LIST_VIEW/MAP_VIEW仅用于聚合累加器底层通过to_sql_type/from_sql_type/need_conversion在 UDF 边界完成 Python 对象与内部 SQL 对象的互转复合类型会递归委托给元素类型。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐PyFlink Table API 数据类型完全指南Data Types 详解与 Python/Pandas 类型映射PyFlink Table API 数据类型完全指南Data Types 详解与 Python/Pandas 类型映射 本文基于 docs/content/d大数据流处理批处理数据工程Daft 类型转换完全指南DataType 与 Python 类型之间的双向映射机制Daft 类型转换完全指南DataType 与 Python 类型之间的双向映射机制 本文围绕 Daft 官方文档 Type Conversions http大数据数据分析数据工程AI 应用PySpark 类型转换完全指南Python 与 Spark SQL 数据类型映射、配置与实战PySpark 类型转换完全指南Python 与 Spark SQL 数据类型映射、配置与实战 本指南以 Apache Spark 仓库中的官方教程 type大数据数据分析批处理流处理机器学习图计算创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
企业数字化 ERP 产品动态
相关推荐
干货实录 | Hypervisor技术在功能安全架构中的应用 本文整理自SASETECH社区技术分享直播《Hypervisor技术在功能安全架构中的应用》,主讲人是小鹏汽车智驾SoC功能安全架构师李智宇老师——深耕汽车电子领域十余年,熟悉行业内主流的汽车软件技术与开发方法。专注功能安全架构设计逾五年,成功主导… · 2026/9/23 10:37:50
OpenStock开源实践:从零搭建个人股票数据分析与可视化系统 最近在复盘自己的数据工作流时,我经常被问到一个问题:个人做股票分析,到底要不要自己搭一套行情数据系统?市面上的炒股软件、免费网页看盘工具一抓一大把,为什么还要费劲自己动手?如果你也有同样的疑问&… · 2026/9/23 10:37:50
PHP-CS-Fixer `indentation_type` 规则详解:统一缩进风格,强制 PSR-2 缩进规范 PHP-CS-Fixer indentation_type 规则详解:统一缩进风格,强制 PSR-2 缩进规范 【免费下载链接】PHP-CS-Fixer A tool to automatically fix PHP Coding Standards issues 项目地址: https://gitcode.com/gh_mirrors/ph/PHP-CS-Fixer
导读
indenta… · 2026/9/23 10:37:37
Word如何单独删除某一页水印?分节符与页眉链接设置详解 1. 为什么“去掉当前页水印”会这么麻烦先说结论:Word里的水印,本质上不是“贴在某一页上的”,而是“整节共享的页眉层内容”。所以如果你直接在文档里删水印,它会整篇消失;你也不可能像选一张图片那样,用鼠… · 2026/9/23 11:22:15
偶滴性能优化保姆级教程:面试答不上来?3招搞定 偶滴性能优化保姆级教程:面试答不上来?3招搞定 面试被问原理答不上来,是不是心里直打鼓?别慌,这份偶滴性能优化保姆级教程,专治各种“卡顿焦虑”。很多开发者以为偶滴只是个小工具,其实它在高并发场景下的瓶颈比想象中更隐蔽。今天我们就用实战数据说… · 2026/9/23 11:22:15
县域商业信息组织:基于相册导航的本地化货源平台设计 1. 项目概述:这不是一个“系统”,而是一套本地化商业信息组织方法“安福相册导航系统”和“安福货源平台源码”这两个词,最近在一些区域性的建站论坛、微信公众号推文和本地服务商报价单里频繁出现。但必须先说清楚:它不是一个像微… · 2026/9/23 11:22:15
单通道EEG睡眠分期实战:从Python清洗到树莓派实时部署 简介:本资源是一项基于单通道脑电信号的自动睡眠分期研究实现,面向生物医学工程、人工智能交叉领域的初学者与进阶学习者,适用于课程设计、毕业设计及科研入门项目。项目复现并改进了TinySleepNet架构,集成双向RNN、GRU与Attentio… · 2026/9/23 11:22:15
城市驾驶系统源码避坑指南:版本升级API重构实战 城市驾驶系统源码避坑指南:版本升级API重构实战 昨天刚把老项目升级到 v2.0,一跑起来,满屏的红字报错。以前用的 drive(city) 接口直接没了,替换成 navigate(location) 还得传一堆新参数。这种“版本升级后… · 2026/9/23 11:22:14
麦克风怎么选?从换能原理到AI协同时代的底层逻辑 经常有朋友私信问我:直播录音买什么麦克风、录书用什么麦、几百块和上千块的差别到底在哪。一开始我都是直接报型号,直到有位做播客的朋友拿着某款热门口碑麦回去,在自家没有做任何声学处理的房间里录了一期节目,人声闷成一团&… · 2026/9/23 11:22:08
3招搞定手机怎么下载微信面试难题实战项目解析 3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29