【问题标题】:how to use sparkSession in dataframe write in pyspark using spark-cassandra-connector如何在数据帧中使用 sparkSession 使用 spark-cassandra-connector 在 pyspark 中写入
【发布时间】:2020-09-30 03:50:14
【问题描述】:

我将 pysparkspark-cassandra-connector_2.11-2.3.0.jar 与 cassandra DB 一起使用。 我正在从一个键空间读取数据帧并写入另一个不同的键空间。这两个键空间有不同的用户名和密码。

我使用以下方法创建了 sparkSession:

spark_session = None

def set_up_spark(sparkconf,config):
    """
    sets up spark configuration and create a session
    :return: None
    """
    try:
        logger.info("spark conf set up Started")
        global spark_session
        spark_conf = SparkConf()
        for key, val in sparkconf.items():
            spark_conf.set(key, val)
        spark_session = SparkSession.builder.config(conf=spark_conf).getOrCreate()
        logger.info("spark conf set up Completed")
    except Exception as e:
        raise e

我使用此 sparkSession 将数据作为数据帧读取为:

table_df = spark_session.read \
            .format("org.apache.spark.sql.cassandra") \
            .options(table=table_name, keyspace=keyspace_name) \
            .load()

我可以使用上述会话读取数据。 spark_session 附加到上述查询。

现在我需要创建另一个会话,因为写入表的凭据不同。 我有写查询:

table_df.write \
            .format("org.apache.spark.sql.cassandra") \
            .options(table=table_name, keyspace=keyspace_name) \
            .mode("append") \
            .save()

我找不到如何在 cassandra 中为上述写入操作附加新的 sparkSession。

如何使用 spark-cassandra-connector 在 pyspark 中附加新的 SparkSession 以进行写入操作?

【问题讨论】:

    标签: python apache-spark pyspark cassandra spark-cassandra-connector


    【解决方案1】:

    您可以简单地将这些信息作为选项传递给特定的readwrite 操作,这包括:spark.cassandra.connection.host

    请注意,您需要将这些选项放入字典中,并传递此字典而不是直接传递,如 documentation 中所述。

    read_options = { "table": "..", "keyspace": "..", 
      "spark.cassandra.connection.host": "IP1", 
      "spark.cassandra.auth.username": "username1", 
      "spark.cassandra.auth.password":"password1"}
    table_df = spark_session.read \
                .format("org.apache.spark.sql.cassandra") \
                .options(**read_options) \
                .load()
    
    write_options = { "table": "..", "keyspace": "..", 
      "spark.cassandra.connection.host": "IP2", 
      "spark.cassandra.auth.username": "username2", 
      "spark.cassandra.auth.password":"password1"}
    table_df.write \
                .format("org.apache.spark.sql.cassandra") \
                .options(**write_options) \
                .mode("append") \
                .save()
    

    【讨论】:

      猜你喜欢
      • 2022-10-02
      • 2019-10-15
      • 2017-02-08
      • 2019-10-31
      • 2015-05-29
      • 2018-12-10
      • 2016-09-02
      • 2018-01-24
      • 2022-07-19
      相关资源
      最近更新 更多