【问题标题】:Flink how to create table with the schema inferred from Avro input dataFlink 如何使用从 Avro 输入数据推断出的模式创建表
【发布时间】:2019-01-28 21:54:28
【问题描述】:

我在 Flink 数据集中加载了一个 Avro 文件:

AvroInputFormat<GenericRecord> test = new AvroInputFormat<GenericRecord>(
        new Path("PathToAvroFile")
        , GenericRecord.class);
DataSet<GenericRecord> DS = env.createInput(test);

usersDS.print();

这是打印 DS 的结果:

{"N_NATIONKEY": 14, "N_NAME": "KENYA", "N_REGIONKEY": 0, "N_COMMENT": " pending excuses haggle furiously deposits. pending, express pinto beans wake fluffily past t"}
{"N_NATIONKEY": 15, "N_NAME": "MOROCCO", "N_REGIONKEY": 0, "N_COMMENT": "rns. blithely bold courts among the closely regular packages use furiously bold platelets?"}
{"N_NATIONKEY": 16, "N_NAME": "MOZAMBIQUE", "N_REGIONKEY": 0, "N_COMMENT": "s. ironic, unusual asymptotes wake blithely r"}
{"N_NATIONKEY": 17, "N_NAME": "PERU", "N_REGIONKEY": 1, "N_COMMENT": "platelets. blithely pending dependencies use fluffily across the even pinto beans. carefully silent accoun"}
{"N_NATIONKEY": 18, "N_NAME": "CHINA", "N_REGIONKEY": 2, "N_COMMENT": "c dependencies. furiously express notornis sleep slyly regular accounts. ideas sleep. depos"}
{"N_NATIONKEY": 19, "N_NAME": "ROMANIA", "N_REGIONKEY": 3, "N_COMMENT": "ular asymptotes are about the furious multipliers. express dependencies nag above the ironically ironic account"}
{"N_NATIONKEY": 20, "N_NAME": "SAUDI ARABIA", "N_REGIONKEY": 4, "N_COMMENT": "ts. silent requests haggle. closely express packages sleep across the blithely"}

现在我想从 DS Dataset 创建一个与 Avro 文件架构完全相同的表,我的意思是列应该是 N_NATIONKEY、N_NAME、N_REGIONKEY 和 N_COMMENT。

我知道使用这条线:

tableEnv.registerDataSet("tbTest", DS, "field1, field2, ...");

我可以创建一个表并设置列,但我希望从数据中自动推断出列。可能吗? 另外,我试过了

tableEnv.registerDataSet("tbTest", DS);

但它会使用架构创建一个表:

root
 |-- f0: GenericType<org.apache.avro.generic.GenericRecord>

【问题讨论】:

    标签: apache-flink flink-sql


    【解决方案1】:

    GenericRecord 是 Table & SQL API 运行时的黑盒,因为字段数量及其数据类型未定义。我建议使用扩展 SpecificRecord 的 Avro 生成的类。 Flink 的类型系统也可以识别这些特定类型,您可以使用适当的数据类型正确处理各个字段。

    或者,您可以实现一个custom UDF,以提取具有适当数据类型getAvroInt(f0, "myField")getAvroString(f0, "myField") 等的字段。

    一些伪代码:

    class AvroStringFieldExtract extends ScalarFunction {
        public String eval(GenericRecord r, String fieldName) {
            return r.get(fieldName).toString();
        }
    }
    
    tableEnv.registerFunction("getAvroFieldString", new AvroStringFieldExtract())
    

    【讨论】:

    • 谢谢,这让我有点困惑。请给我一个使用自定义 UDF 的简单示例吗?
    • 我在答案中添加了一个小例子。希望这会有所帮助。
    • 谢谢。我应该为每个 Avro 字段扩展一个 ScalarFunction 类吗?
    • 否,仅适用于您要提取的数据类型。 fieldName 参数使函数具有通用性。
    猜你喜欢
    • 2019-08-19
    • 1970-01-01
    • 1970-01-01
    • 2021-02-26
    • 1970-01-01
    • 2019-02-08
    • 1970-01-01
    • 2014-09-24
    • 2022-11-11
    相关资源
    最近更新 更多