【发布时间】: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