相关疑难解决方法(0)

使用模式将带有Spark的AVRO消息转换为DataFrame

有没有使用模式转换方式从消息与到?用户记录的模式文件:

{
  "fields": [
    { "name": "firstName", "type": "string" },
    { "name": "lastName", "type": "string" }
  ],
  "name": "user",
  "type": "record"
}
Run Code Online (Sandbox Code Playgroud)

来自SqlNetworkWordCount示例和Kafka,Spark和Avro的代码片段- 第3部分,生成和使用Avro消息来读取消息.

object Injection {
  val parser = new Schema.Parser()
  val schema = parser.parse(getClass.getResourceAsStream("/user_schema.json"))
  val injection: Injection[GenericRecord, Array[Byte]] = GenericAvroCodecs.toBinary(schema)
}

...

messages.foreachRDD((rdd: RDD[(String, Array[Byte])]) => {
  val sqlContext = SQLContextSingleton.getInstance(rdd.sparkContext)
  import sqlContext.implicits._

  val df = rdd.map(message => Injection.injection.invert(message._2).get)
    .map(record => User(record.get("firstName").toString, records.get("lastName").toString)).toDF()

  df.show() …
Run Code Online (Sandbox Code Playgroud)

scala avro apache-kafka apache-spark spark-streaming

13
推荐指数
1
解决办法
1万
查看次数

如何在Spark中引入一行模式?

在Row Java API中有一个row.schema(),但是没有row.set(StructType模式).

我也试过RowFactorie.create(objets),但我不知道如何继续

更新:

问题是当我修改示例中的工作者的结构时如何生成新的数据帧

DataFrame sentenceData = jsql.createDataFrame(jrdd, schema);
List<Row> resultRows2 = sentenceData.toJavaRDD()
            .map(new MyFunction<Row, Row>(parameters) {
            /** my map function **// 

                public Row call(Row row) {

                 // I want to change Row definition adding new columns
                    Row newRow = functionAddnewNewColumns (row);
                    StructType newSchema = functionGetNewSchema (row.schema);

                    // Here I want to insert the structure 

                    //
                    return newRow
                    }

                }

        }).collect();


JavaRDD<Row> jrdd = jsc.parallelize(resultRows);

// Here is the problema  I don't know how to get the new …
Run Code Online (Sandbox Code Playgroud)

apache-spark

7
推荐指数
2
解决办法
1万
查看次数

Avro Schema引发StructType

这实际上与我之前的问题相同,但使用Avro而不是JSON作为数据格式.

我正在使用Spark数据帧,它可以从几个不同的模式版本之一加载数据:

// Version One
{"namespace": "com.example.avro",
 "type": "record",
 "name": "MeObject",
 "fields": [
     {"name": "A", "type": ["null", "int"], "default": null}
 ]
}

// Version Two
{"namespace": "com.example.avro",
 "type": "record",
 "name": "MeObject",
 "fields": [
     {"name": "A", "type": ["null", "int"], "default": null},
     {"name": "B", "type": ["null", "int"], "default": null}
 ]
}
Run Code Online (Sandbox Code Playgroud)

我正在使用Spark Avro加载数据.

DataFrame df = context.read()
  .format("com.databricks.spark.avro")
  .load("path/to/avro/file");
Run Code Online (Sandbox Code Playgroud)

可以是Version One文件或Version Two文件.但是我希望能够以相同的方式处理它,将未知值设置为"null".我之前的问题中的建议是设置模式,但是我不想重复自己在.avro文件和火花StructType和朋友中编写模式.如何将avro架构(文本文件或生成的MeObject.getClassSchema())转换为火花StructType?

Spark Avro有一个SchemaConverters,但它都是私有的,并返回一些奇怪的内部对象.

java avro apache-spark apache-spark-sql

4
推荐指数
1
解决办法
9011
查看次数

如何将RDD [GenericRecord]转换为scala中的dataframe?

我从Avaf(序列化器和反序列化器)获得了kafka主题的推文.然后我创建了一个Spark消费者,它在RDD [GenericRecord]的Dstream中提取推文.现在我想将每个rdd转换为数据帧,以通过SQL分析这些推文.有什么解决方案可以将RDD [GenericRecord]转换为数据帧吗?

scala avro apache-spark spark-dataframe

1
推荐指数
1
解决办法
3377
查看次数