【发布时间】:2019-04-23 10:57:26
【问题描述】:
我是 spark 新手,正在尝试探索 Spark 结构化流。我将使用来自 Kafka(嵌套 JSON)的消息,根据 JSON 属性上的某些条件过滤这些消息。然后应将满足过滤器的每条消息推送到 Cassandra。
我已阅读有关 spark Cassandra 连接器的文档 https://spark.apache.org/docs/2.2.0/structured-streaming-kafka-integration.html
Dataset<Row> df = spark
.readStream()
.format("kafka")
.option("kafka.bootstrap.servers", "host1:port1,host2:port2")
.option("subscribe", "topic1")
.load()
df.selectExpr("CAST(value AS STRING)")
我只需要这个嵌套 JSON 中存在的众多属性中的几个。如何在其之上应用架构,以便可以使用 sparkSQL 进行过滤?
对于示例 JSON,我需要为玩频率总和超过 5 的玩家坚持姓名、年龄、经验、爱好名称、爱好经验。
{
"name": "Tom",
"age": "24",
"gender": "male",
"hobbies": [{
"name": "Tennis",
"experience": 5,
"places": [{
"city": "London",
"frequency": 4
}, {
"city": "Sydney",
"frequency": 3
}]
}]
}
我对 Spark 比较陌生,如有重复请见谅。另外,我正在寻找 JAVA 中的解决方案。
【问题讨论】:
标签: apache-spark spark-structured-streaming