代码之家  ›  专栏  ›  技术社区  ›  Soheil Pourbafrani

Flink如何使用从Avro输入数据推断的模式创建表

  •  0
  • Soheil Pourbafrani  · 技术社区  · 7 年前

    我已将Avro文件加载到Flink数据集中:

    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数据集中创建一个表,它的模式和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>
    
    1 回复  |  直到 7 年前
        1
  •  1
  •   twalthr    7 年前

    GenericRecord 是桌子上的黑盒子&SQL API运行时字段数及其数据类型未定义。我建议使用Avro生成的类来扩展 SpecificRecord . 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())
    
    推荐文章