【问题标题】:Add schema to a Dataset[Row] in Java在 Java 中将模式添加到 Dataset[Row]
【发布时间】: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


    【解决方案1】:

    您可以像这样指定您的架构:

    import org.apache.spark.sql.types.{DataTypes, StructField, StructType};
    
    StructType schema = DataTypes.createStructType(new StructField[] {
        DataTypes.createStructField("name",  DataTypes.StringType, true),
        DataTypes.createStructField("age", DataTypes.StringType, true),
        DataTypes.createStructField("gender", DataTypes.StringType, true),
        DataTypes.createStructField("hobbies", DataTypes.createStructType(new StructField[] {
            DataTypes.createStructField("name", DataTypes.StringType, true),
            DataTypes.createStructField("experience", DataTypes.IntegerType, true),
            DataTypes.createStructField("places", DataTypes.createStructType(new StructField[] {
                DataTypes.createStructField("city", DataTypes.StringType, true),
                DataTypes.createStructField("frequency", DataTypes.IntegerType, true)
            }), true)
        }), true)
    });
    

    然后根据需要使用架构创建数据框:

    import org.apache.spark.sql.functions.{col, from_json};
    
    df.select(from_json(col("value"), schema).as("data"))
      .select(
        col("data.name").as("name"),
        col("data.hobbies.name").as("hobbies_name"))
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2020-10-21
      • 2020-05-11
      • 2018-10-02
      • 2018-10-30
      • 2018-10-01
      • 1970-01-01
      • 2017-11-17
      • 2018-09-11
      相关资源
      最近更新 更多