Flink 架构与表架构

Nak*_*euh 1 java apache-flink flink-sql

我正在使用 Flink SQL API,我在所有“模式”类型之间有点迷失:TableSchema、Schema(来自org.apache.flink.table.descriptors.Schema)和TypeInformation。

ATableSchema可以从 a 创建TypeInformation,aTypeInformation可以从 a 创建TableSchema,aSchema可以从 a 创建TableSchema

但看起来 aSchema无法转换回TypeInformationor TableSchema(?)

为什么有 3 种不同类型的对象来存储同一种信息?

例如,假设我有一个来自 Avro 架构文件的字符串架构,并且我想向其中添加一个新字段。为此,我找到的唯一解决方案是:

String mySchemaRaw = ...;
TypeInformation<Row> typeInfo = AvroSchemaConverter.convertToTypeInfo(mySchemaRaw);
Schema newSchema = new Schema().schema(TableSchema.fromTypeInfo(typeInfo));
newSchema = newSchema.field("nexField",...);


// Need the newSchema as a TableSchema 
Run Code Online (Sandbox Code Playgroud)

这是使用这些对象的正常方式吗?(我觉得很奇怪)

twa*_*thr 5

TypeInformation并TableSchema解决不同的事情。TypeInformation是如何将记录类(例如行或 POJO)从一个操作员传送到另一个操作员的物理信息。

TableSchema描述独立于底层每记录类型的表模式。它类似于CREATE TABLE name (a INT, b BIGINT)DDL 语句的架构部分。在 SQL 中,也没有定义像CREATE TABLE name ROW(a INT, B BIGINT). 但架构和行类型确实相关,这就是提供转换器方法的原因。PRIMARY KEY一旦引入诸如此类的概念,差异就会变得更大。

Schema是当前指定非 SQL 概念(例如时间属性和字段映射)的方法。