文章 · 2022-02-24

Cassandra 数据清理实战

我们设计了一套五步数据清理方案:

  1. 导出全部商品ID:从Cassandra中导出 products 表的所有商品ID,生成完整的ID列表。
  2. 获取有效商品ID列表:从业务系统获取当前有效(在售、有效)商品ID列表,作为权威数据源。
  3. 筛选无效商品ID:将全量ID列表与有效ID列表比对,找出无效的商品ID—即在Cassandra中存在但已失效的ID。
  4. 批量删除无效数据:通过脚本批量删除这些无效商品ID对应的数据记录。
  5. 验证清理效果:删除完成后,重新检查数据量或抽检商品ID,确保无效数据已被移除。

工作流程为:数据导出 → 有效ID获取 → 无效ID筛选 → 批量删除 → 结果验证。

数据导出:导出Cassandra商品ID

第一步是获取当前Cassandra中存储的所有商品ID。假设商品数据表名为 products,我们只需导出该表的主键 id 列。

使用Cassandra自带的 cqlsh 工具的 COPY 命令导出数据。为了处理数千万行的数据量,导出前应适当提高超时时间:

$ cqlsh --request-timeout=600000 <cassandra_host>
cqlsh> USE product_keyspace;
cqlsh:product_keyspace> COPY products(id) TO '/tmp/products_ids.csv' WITH HEADER = false;

该命令将 products 表的 id 列导出到服务器临时目录的CSV文件(不包含表头)。导出数千万行ID通常耗时十余分钟。导出完成后,通过 scp 将文件拷贝到本地处理:

$ scp user@<cassandra_server>:/tmp/products_ids.csv ./products_ids.csv

现在本地有了 products_ids.csv 文件,包含所有商品ID。

获取有效商品ID列表

接下来,从业务系统获取当前有效商品ID列表,作为判定无效数据的权威来源。通常来自电商商品服务或类似的权威系统。例如,若内部API返回JSON格式的商品ID数组,可用 curl 调用并保存结果:

$ curl "http://internal.api.company/active_products" -H "Content-Type: application/json" -d '{"pageSize":10000,"currentPage":1}' -o active_products.json

API返回结果可能需要简单处理—提取商品ID字段并保存到 valid_ids.txt,每行一个ID。

筛选无效商品ID

现在对比Cassandra导出的全量ID列表和有效ID列表。无效ID即在Cassandra中存在但不在有效列表中的ID。

使用此Python脚本进行筛选:

# load valid ids into a set for fast lookup
valid_ids = set()
with open('valid_ids.txt', 'r') as f:
    for line in f:
        vid = line.strip()
        if vid:
            valid_ids.add(vid)

# iterate over all exported IDs and filter
count = 0
with open('products_ids.csv', 'r') as f_all, open('garbage_ids.txt', 'w') as f_out:
    for line in f_all:
        pid = line.strip()
        if not pid:
            continue
        if pid not in valid_ids:
            f_out.write(pid + "\n")
            count += 1

f_out.close()
print("无效商品ID数量:", count)

脚本逐行读取 products_ids.csv,检查每个ID是否在 valid_ids 集合中。不在集合中的ID被写入 garbage_ids.txt。最后打印无效ID的总数。

注: 处理大文件时,必须逐行读取以避免内存溢出。此例假定有效ID数量足以装入内存。

删除无效数据

有了待删除ID列表后,使用Python Cassandra驱动批量删除:

from cassandra.cluster import Cluster

# 连接Cassandra集群
cluster = Cluster(['<cassandra_host_ip>'], port=9042)
session = cluster.connect('product_keyspace')

# 准备参数化删除语句
delete_stmt = session.prepare("DELETE FROM products WHERE id = ?")

# 逐行读取无效ID文件,执行删除
count = 0
with open('garbage_ids.txt', 'r') as f:
    for line in f:
        pid = line.strip()
        if not pid:
            continue
        session.execute(delete_stmt, [pid])
        count += 1

print("删除完成,删除总条数:", count)

脚本读取 garbage_ids.txt 中的每个ID并执行 DELETE 语句。预编译语句可提高效率。删除数百万条记录需要时间—应在业务低峰期执行。

注意: Cassandra的 DELETE 操作不会立即移除数据,而是标记为"墓碑"等待后续压缩清理。大规模删除会积累墓碑,可能影响后续读性能。需要监控集群状况,适时进行压缩等维护。

验证清理结果

删除完成后,通过两种方式验证成功:

验证无误后,清理工作完成。

实践总结

本次清理实践的主要经验:

通过此过程,我们成功清理了Cassandra中大量无效商品数据,减少搜索系统的负担,同时保护了线上数据准确性。

© 2026 Yuxu Ge ·