【问题标题】:Spark dataframe from dynamic schema Or filter out rows that do not satisfy the schema从动态模式中激发数据帧或过滤掉不满足模式的行
【发布时间】:2018-01-15 13:34:49
【问题描述】:

我编写了一个 kafka 生产者,它跟踪日志文件(格式:csv)的内容。kafka 消费者是一个创建 JavaDStream 的流应用程序。 使用 forEachRDD 方法,我将文件的每一行拆分为分隔符“,”并创建 Row 对象。我指定了具有 7 列的模式。 然后我使用 JavaRDD 和模式创建数据框。 但这里的问题是,日志文件中的所有行都没有相同的列数。 因此,有没有办法过滤掉这些不满足模式的行或根据行内容动态创建模式? 以下是部分代码:

JavaDStream<String> msgDataStream =directKafkaStream.map(new Function<Tuple2<String, String>, String>() {
            @Override
            public String call(Tuple2<String, String> tuple2) {
                return tuple2._2();
            }
        });
msgDataStream.foreachRDD(new VoidFunction<JavaRDD<String>>() {
            @Override
            public void call(JavaRDD<String> rdd) {
                JavaRDD<Row> rowRDD = rdd.map(new Function<String, Row>() {
                    @Override
                    public Row call(String msg) {
                        String[] splitMsg=msg.split(",");
                        Object[] vals = new Object[splitMsg.length];
                        for(int i=0;i<splitMsg.length;i++)
                        {
                            vals[i]=splitMsg[i].replace("\"","").trim();
                        }

                        Row row = RowFactory.create(vals);
                        return row;
                    }
                });

                //Create Schema
                StructType schema = DataTypes.createStructType(new StructField[] {
                        DataTypes.createStructField("timeIpReq", DataTypes.StringType, true),DataTypes.createStructField("SrcMac", DataTypes.StringType, true),
                        DataTypes.createStructField("Proto", DataTypes.StringType, true),DataTypes.createStructField("ACK", DataTypes.StringType, true),
                        DataTypes.createStructField("srcDst", DataTypes.StringType, true),DataTypes.createStructField("NATSrcDst", DataTypes.StringType, true),
                        DataTypes.createStructField("len", DataTypes.StringType, true)});

                //Get Spark 2.0 session


                Dataset<Row> msgDataFrame = session.createDataFrame(rowRDD, schema);

【问题讨论】:

  • 你需要使用Java吗?我可以在 Scala 中为您提供解决方案。会有帮助吗?
  • 是的..我使用 Java 是因为我更熟悉它。但是如果您也可以在 Scala 中提供解决方案,那将有很大帮助!谢谢!

标签: apache-spark apache-kafka apache-spark-sql spark-dataframe spark-streaming


【解决方案1】:

删除与预期架构不匹配的行的一种简单方法是将flatMapOption 类型一起使用,此外,如果您的目标是构建DataFrame,我们使用相同的flatMap 步骤来应用数据的架构。这在 Scala 中通过使用case classes 来促进。

// Create Schema
case class NetInfo(timeIpReq: String, srcMac: String, proto: String, ack: String, srcDst: String, natSrcDst: String, len: String)

val netInfoStream = msgDataStream.flatMap{msg => 
  val parts = msg.split(",")
  if (parts.size == 7) {  //filter out messages with unmatching set of fields
    val Array(time, src, proto, ack, srcDst, natSrcDst, len) = parts // use a extractor to get the different parts in variables
    Some(NetInfo(time, src, proto, ack, srcDst, natSrcDst, len)) // return a valid record
  } else {
    None  // We don't have a valid. Return None
  }
}

netInfoStream.foreachRDD{rdd =>
    import sparkSession.implicits._ 
    val df = rdd.toDF() // DataFrame transformation is possible on RDDs with a schema (based on a case class)
    // do stuff with the dataframe
}

关于:

日志文件中的所有行的列数都不相同。

假设它们都代表相同类型的数据,但可能缺少某些列,正确的策略是过滤掉不完整的数据(如此处示例)或在定义的架构中使用可选值(如果存在确定性)知道缺少哪些字段的方法。应该向生成数据的上游应用程序提出此要求。在 CSV 中用空逗号序列表示缺失值是很常见的(例如field0,,field2,,,field5

处理每行差异的动态架构没有意义,因为无法将其应用于由具有不同架构的行组成的 DataFrame

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-11-12
    • 2022-12-21
    • 2016-12-20
    • 2021-03-21
    相关资源
    最近更新 更多