在Databricks中使用Python Cassandra Driver查询Cassandra表的问题与解决
Databricks中Cassandra查询优化问题排查与解决
问题背景
在Databricks中尝试优化Cassandra表查询,采用逐行查询合并结果的方案,使用Python Cassandra Driver实现时遇到两处问题,最终通过Spark DataFrame解决。
阶段一:NoneType无法迭代问题
初始代码
from cassandra.cluster import Cluster from cassandra.auth import PlainTextAuthProvider import pandas as pd def init_cassandra_session(endpoints, keyspace, username, password, port=9042): auth_provider = PlainTextAuthProvider(username, password) cluster = Cluster(endpoints, port=port, auth_provider=auth_provider) cassandra_session = cluster.connect(keyspace, wait_for_all_pools=False) return cassandra_session def get_rdd_values(rows): out_df = None cassandra_session = init_cassandra_session(host, keyspace, username, password) for row in rows: device_id = row.__getitem__('device_id') timestamp = row.__getitem__('timestamp') category = row.__getitem__('category') query = ''' select * from headcounter_category_h_aggr where device_id = '%s' and timestamp = '&s' and category = '%s' ''' result_query = cassandra_session.execute(query, [device_id, timestamp, category]) if out_df is None: out_df = result_query else: out_df = out_df.append(result_query) return out_df columns = ['device_id', 'timestamp', 'category'] data = [['SIMUL_TEST03', '2020-12-23 11:00:00', 'PERSON'], ['SIMUL_TEST03', '2020-12-23 12:00:00', 'PERSON']] pdf = pd.DataFrame(data, columns=columns) dfFromData1 = spark.createDataFrame(pdf) rdd_values = dfFromData1.rdd.mapPartitions(get_rdd_values) rdd_values.collect()
问题根源
get_rdd_values中初始out_df设为None,若无查询结果会直接返回None,导致RDD无法迭代- 查询语句存在语法错误:
timestamp = '&s'应为'%s',且字符串拼接+参数传递的写法易引发问题
阶段二:查询结果重复问题
修改后的get_rdd_values函数
def get_rdd_values(rows): out_df = [] cassandra_session = init_cassandra_session(host, keyspace, username, password) for row in rows: device_id = row.__getitem__('device_id') timestamp = row.__getitem__('timestamp') category = row.__getitem__('category') query = f"select * from headcounter_category_h_aggr where device_id = '{device_id}' and timestamp = '{timestamp}' and category = '{category}'" result_query = cassandra_session.execute(query) if len(out_df)== 0: out_df = result_query else: out_df = out_df.append(result_query) return out_df
问题根源
- 初始
out_df为列表,但后续直接赋值为Cassandra查询返回的ResultSet对象 ResultSet.append()操作不会返回新对象,直接拼接会导致结果结构混乱,引发重复
最终解决方案
将RDD转换为Spark DataFrame,利用Spark的分布式处理特性规范结果结构,彻底解决重复问题:
dfFromRDD = spark.createDataFrame(rdd_values, schema = schema)
额外优化建议
- 避免在
mapPartitions中循环逐行查询,改为批量查询减少Cassandra连接开销 - 优先使用Cassandra官方Spark Connector替代Python Driver,适配Spark生态且性能更优
- 始终采用参数化查询,规避SQL注入风险同时避免语法错误
内容的提问来源于stack exchange,提问作者Gabriele Sciurti
相关产品推荐
相关产品推荐

