【问题标题】:Spark sql create dataset from stringSpark sql从字符串创建数据集
【发布时间】:2020-07-04 03:15:00
【问题描述】:

我有一个字符串,我想从中创建一个数据集,该字符串是 \n 分隔的行和 \t 分隔的字段:
8 "SOMETHING" 15236236 "2" "SOMETHING" "SOMETHHING"
所以我将字符串拆分为\n 并从中创建一个List<String>,然后我使用JavaSparkContext 实例创建一个JavaRDD,然后我尝试使用sqlContet 方法createDataset 创建我的数据集。
这编译得很好,如果我在 loadDataset() 方法的 return 语句处设置一个断点,我会看到 settingsDataset 数据集,它只会在代码调用第一个操作后中断。

我试图实现这一点的方式是:

private Dataset<Row> loadDataset(){
    InputStream in;
    Dataset<Row> settingsDataset = null;
    try {
      JavaSparkContext jsc = new JavaSparkConte xt(session.sparkContext());
      in = getClass().getResourceAsStream("filename.tsv");
      String settingsFileAsString = IOUtils.toString(in, Charsets.UTF_8);
      List<String> settingsFileAsList = Arrays.asList(settingsFileAsString.split("\n"));
      Encoder<Row> encoder = RowEncoder.apply(getSchema());
      JavaRDD settingsFileAsRDD = jsc.parallelize(settingsFileAsList);
      settingsDataset = session.sqlContext().createDataset(settingsFileAsRDD.rdd(), encoder).toDF();
    } catch (Exception e) {
      e.printStackTrace();
    }
    return settingsDataset;
 }

  private org.apache.spark.sql.types.StructType getSchema() {
    return DataTypes.createStructType(new StructField[]{
        DataTypes.createStructField("f_1", DataTypes.StringType, true),
        DataTypes.createStructField("f_2", DataTypes.StringType, true),
        DataTypes.createStructField("f_3", DataTypes.StringType, true),
        DataTypes.createStructField("f_4", DataTypes.StringType, true),
        DataTypes.createStructField("f_5", DataTypes.StringType, true),
        DataTypes.createStructField("f_6", DataTypes.StringType, true)
    });
  }

问题是 DAG 不会被创建,代码中断并出现以下异常: ! java.lang.ClassCastException: java.lang.String cannot be cast to org.apache.spark.sql.Row

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    实际上JavaRDD settingsFileAsRDD = jsc.parallelize(settingsFileAsList);JavaRDD&lt;String&gt; 但它应该是JavaRDD&lt;Row&gt;。您应该用\t 拆分这些“行”,并使用RowFactory.create(s.split("\t")) 创建新的Row。请参见下面的示例:

    SparkSession spark = SparkSession.builder().master("local").getOrCreate();
    JavaSparkContext jsc = new JavaSparkContext(spark.sparkContext());
    String settingsFileAsString = "1\t2\t3\t4\t5\t6\n7\t8\t9\t10\t11\t12";
    List<String> settingsFileAsList = Arrays.asList(settingsFileAsString.split("\n"));
    Encoder<Row> encoder = RowEncoder.apply(getSchema());
    JavaRDD<Row> settingsFileAsRDD = jsc.parallelize(settingsFileAsList).map(s->RowFactory.create(s.split("\t")));
    Dataset<Row> settingsDataset = spark.createDataset(settingsFileAsRDD.rdd(), encoder).toDF();
    settingsDataset.show();
    

    结果:

    +---+---+---+---+---+---+
    |f_1|f_2|f_3|f_4|f_5|f_6|
    +---+---+---+---+---+---+
    |  1|  2|  3|  4|  5|  6|
    |  7|  8|  9| 10| 11| 12|
    +---+---+---+---+---+---+
    

    【讨论】:

      猜你喜欢
      • 2020-02-29
      • 1970-01-01
      • 2018-10-23
      • 1970-01-01
      • 1970-01-01
      • 2017-02-19
      • 2021-04-01
      • 2016-10-14
      • 2012-11-22
      相关资源
      最近更新 更多