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

如何在Cassandra中实现跨多表查询最高频次航空公司?

解决Cassandra多表统计航空公司频次的思路

针对Cassandra不支持跨表查询的问题,这里提供几个可行的解决方向:

1. 提前设计聚合表(推荐,符合Cassandra设计理念)

Cassandra是查询优先的数据库,最佳实践是在数据导入阶段就构建适配统计需求的表结构:

  • 创建一张专门用于统计航空公司频次的表:
CREATE TABLE airline_frequency (
    airline_code TEXT PRIMARY KEY,
    airline_name TEXT,
    total_count BIGINT DEFAULT 0
);
  • 导入各年份CSV数据到对应表的同时,同步更新这张聚合表:
    用UPDATE语句累加计数,批量处理时每读取一条记录执行:
    UPDATE airline_frequency 
    SET total_count = total_count + 1 
    WHERE airline_code = ?;
    
  • 后续直接查询这张表就能快速得到结果:
    SELECT airline_name, total_count 
    FROM airline_frequency 
    ORDER BY total_count DESC LIMIT 1;
    

2. 用大数据工具做跨表聚合

如果已完成数据导入且不想修改现有表结构,可以借助Spark/Flink这类工具连接Cassandra,拉取多表数据后做聚合:

  • 以Spark为例,使用Spark Cassandra Connector读取所有年份表并合并统计:
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .appName("AirlineTopFrequency") \
    .config("spark.cassandra.connection.host", "你的Cassandra地址") \
    .config("spark.jars.packages", "com.datastax.spark:spark-cassandra-connector_2.12:3.4.1") \
    .getOrCreate()

# 列出所有年份表
table_list = ["airline_data_2018", "airline_data_2019", "airline_data_2020"]

# 读取每张表并合并
combined_data = None
for table in table_list:
    df = spark.read.format("org.apache.spark.sql.cassandra") \
        .options(table=table, keyspace="你的keyspace名称") \
        .load() \
        .select("airline_code", "airline_name")
    if combined_data is None:
        combined_data = df
    else:
        combined_data = combined_data.union(df)

# 统计频次并取最高
top_airline = combined_data.groupBy("airline_code", "airline_name") \
    .count() \
    .orderBy("count", ascending=False) \
    .first()

print(f"频次最高的航空公司:{top_airline['airline_name']},总次数:{top_airline['count']}")

3. 应用层循环查询合并(仅适合小数据量)

如果数据集规模不大,可以在应用程序中循环查询每张年份表,在内存中累加统计:

from cassandra.cluster import Cluster

cluster = Cluster(["你的Cassandra地址"])
session = cluster.connect("你的keyspace名称")

table_list = ["airline_data_2018", "airline_data_2019", "airline_data_2020"]
freq_dict = {}

for table in table_list:
    # 分页查询避免内存溢出
    rows = session.execute(f"SELECT airline_code, airline_name FROM {table}", fetch_size=1000)
    for row in rows:
        key = (row.airline_code, row.airline_name)
        freq_dict[key] = freq_dict.get(key, 0) + 1

# 获取频次最高的条目
top_entry = max(freq_dict.items(), key=lambda x: x[1])
print(f"频次最高的航空公司:{top_entry[0][1]},总次数:{top_entry[1]}")

内容的提问来源于stack exchange,提问作者james kam

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 04:20:17