You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.13 00:25:19