【问题标题】:PySpark Structured Streaming data writing into Cassandra not populating dataPySpark结构化流数据写入Cassandra不填充数据
【发布时间】:2020-02-06 17:13:08
【问题描述】:

我想将 Spark 结构化流数据写入 cassandra。我的 spark 版本是 2.4.0。

我来自 Kafka 的输入源是 JSON,因此在写入控制台时可以,但是当我在 cqlsh Cassandra 中查询时,表中没有附加记录。你能告诉我有什么问题吗?

schema = StructType() \
            .add("humidity", IntegerType(), True) \
            .add("time", TimestampType(), True) \
            .add("temperature", IntegerType(), True) \
            .add("ph", IntegerType(), True) \
            .add("sensor", StringType(), True) \
            .add("id", StringType(), True)

def writeToCassandra(writeDF, epochId):
    writeDF.write \
        .format("org.apache.spark.sql.cassandra") \
        .mode('append') \
        .options("spark.cassandra.connection.host", "cassnode1, cassnode2") \
        .options(table="sensor", keyspace="sensordb") \
        .save()

# Load json format to dataframe
df = spark \
      .readStream \
      .format("kafka") \
      .option("kafka.bootstrap.servers", "kafkanode") \
      .option("subscribe", "iot-data-sensor") \
      .load() \
      .select([
            get_json_object(col("value").cast("string"), "$.{}".format(c)).alias(c)
            for c in ["humidity", "time", "temperature", "ph", "sensor", "id"]])

df.writeStream \
    .foreachBatch(writeToCassandra) \
    .outputMode("update") \
    .start()

【问题讨论】:

    标签: apache-spark cassandra pyspark spark-structured-streaming spark-cassandra-connector


    【解决方案1】:

    我在 pyspark 中遇到了同样的问题。试试下面的步骤

    1. 首先,验证它是否连接到 cassandra。您可以指向一个不可用的表,看看它是否因为“找不到表”而失败

    2. 如下尝试 writeStream(在调用 cassandra 更新之前包括触发器和输出模式)

    df.writeStream \ .trigger(processingTime="10 seconds") \ .outputMode("update") \ .foreachBatch(writeToCassandra) \

    【讨论】:

      猜你喜欢
      • 2019-11-29
      • 2020-06-16
      • 2018-10-06
      • 1970-01-01
      • 1970-01-01
      • 2020-01-27
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多