【发布时间】:2020-09-30 03:50:14
【问题描述】:
我将 pyspark 和 spark-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