【问题标题】:How to extract DB name and table name from a comma separated string of Dataframe column to two column如何从逗号分隔的Dataframe列字符串中提取数据库名称和表名到两列
【发布时间】:2021-02-23 21:59:28
【问题描述】:

我有一个数据框列“table_name”,它的字符串值低于

tradingpartner.parent_supplier,lookup.store,lab_promo_invoice.tl_cc_mbr_prc_wkly_inv,lab_promo_invoice.mpp_club_card_promotion_funding_view,lab_promo_invoice.supplier_sale_apportionment_cc,tradingpartner.supplier,stores.rpm_zone_location_mapping,lookup.calendar

如何从上述字符串中提取数据库名和表名,并将其存储为一列中的数据库名和另一列中的表名。

我想要如下输出

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    您可以使用正则表达式来提取 DBName 和 Table:

    val result = df.select(
        col("table_name"),
        regexp_replace(col("table_name"), "\\.[^,]+(,|$)", "$1").as("DBName"),
        regexp_replace(col("table_name"), "(^|,)[^,]+\\.", "$1").as("Table")
    )
    
    result.show(false)
    +---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+--------------------------------------------------------------------------------------------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------+
    |table_name                                                                                                                                                                                                                                                           |DBName                                                                                                  |Table                                                                                                                                                       |
    +---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+--------------------------------------------------------------------------------------------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------+
    |tradingpartner.parent_supplier,lookup.store,lab_promo_invoice.tl_cc_mbr_prc_wkly_inv,lab_promo_invoice.mpp_club_card_promotion_funding_view,lab_promo_invoice.supplier_sale_apportionment_cc,tradingpartner.supplier,stores.rpm_zone_location_mapping,lookup.calendar|tradingpartner,lookup,lab_promo_invoice,lab_promo_invoice,lab_promo_invoice,tradingpartner,stores,lookup|parent_supplier,store,tl_cc_mbr_prc_wkly_inv,mpp_club_card_promotion_funding_view,supplier_sale_apportionment_cc,supplier,rpm_zone_location_mapping,calendar|
    +---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+--------------------------------------------------------------------------------------------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------+
    

    【讨论】:

    • 非常感谢@mck 的快速响应。我想将输出数据库名称和表名称保存在同一记录中的逗号分隔值中,而不是多个记录或行
    • 编辑了问题。我将输出添加为图像。我希望输出为单个记录中的 table_name DBname Tablename
    • 感谢您的回复。现在,如果我运行上面的 command.org.apache.spark.sql.catalyst.parser.ParseException: extraneous input '>' expecting {'(', 'SELECT'..== SQL == array_join( transform(table_name2,x-> split(x,'\\.')[0]),',') as dbname
    • 你能帮我解决上述错误吗?
    • scala> val result = df.select(col("table_name"), split(col("table_name"), ",").as("table_name2")).selectExpr("table_name ","array_join(transform(table_name2,x -> split(x,'\\\\.')[0]),',') as dbname","array_join(transform(table_name2, x -> split(x , '\\\\.')[1]), ',') as Table") org.apache.spark.sql.catalyst.parser.ParseException: 外部输入 '>' 期望 {'(', 'SELECT' , 'FROM', 'ADD', array_join(transform(table_name2,x -> split(x,'\\.')[0]),',') as dbname ----------- -------------------------^^^ 在 org.apache.spark.sql.catalyst.parser.ParseException.withCommand(ParseDriver.scala:239 )
    【解决方案2】:

    对不起,我用 Java 为您的要求编写了 UDF,但我认为它很容易转换为 Scala。

    //Spark >= 2.3
    UserDefinedFunction splitTableNameUDF = udf(transformTables(), getTableName());
    df.withColumn("table_name_new", splitTableNameUDF.apply(col("table_name")))
        .select("table_name", "table_name_new.DBName", "table_name_new.Table");
    
    
    //Spark < 2.3
    
    sqlContext.udf().register("splitTableNameUDF", transformTables(), getTableName());
    df.withColumn("table_name_new", callUDF("splitTableNameUDF", col("table_name")))
        .select("table_name", "table_name_new.DBName", "table_name_new.Table");
    
    //schema
    
    public static StructType getTableName() {
        List<StructField> inputFields = new ArrayList<>();
        inputFields.add(DataTypes.createStructField("DBName", DataTypes.StringType, true));
        inputFields.add(DataTypes.createStructField("Table", DataTypes.StringType, true));
        return DataTypes.createStructType(inputFields);
    }
    
    
    //UDF
    
    public static UDF1<String, Row> transformTables() {
        return (row) -> {
            String[] dbtableNames = row.split(",");
            List<String> dbNames = new ArrayList<>();
            List<String> tableNames = new ArrayList<>();
            for (String dbtableName : dbtableNames) {
                String[] split = dbtableName.split(".");
                dbNames.add(split[0]);
                tableNames.add(split[1]);
            }
            return RowFactory.create(String.join(",", dbNames), String.join(",", tableNames));
        };
    }
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-08-21
      • 2020-01-21
      • 2012-08-22
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多