如何用Ecto按home_id和path分组获取前3条结果?
问题描述
背景
我有两个查询各自会返回海量数据,无法使用Repo.all,因为这会将数据加载到内存中,很快就会耗尽内存。因此尝试将尽可能多的计算逻辑交给pSQL数据库处理。
现有两个子查询:
统计数量的子查询
all_counts = table_A |> join(:left, [item_A], item_B in table_B, on: item_A.home_id == item_B.home_id and item_A.path == item_B.path ) |> select([unfiltered_item, filtered_item], %{ path: item_A.path, item_fruits_count: coalesce(item_A.fruits, 0), item_veggies_count: coalesce(item_B.veggies, 0), dataset_id: item_A.home_id }) |> subquery()
文件信息子查询
file_info = table_C |> join(:inner, [item], file in table_D, on: item.id == file.item_id and not file.deleted ) |> select([item, file], %{ item_id: item.id, home_id: item.home_id, path: item.path, photo_key: file.photo_key }) |> subquery()
当前问题
合并两个查询的初始方案会返回海量数据,导致内存耗尽:
result = all_counts |> join(:inner, [c], f in ^file_info, on: c.home_id == f.home_id and c.path == f.path) |> select([c, f], %{ item_id: f.item_id, home_id: f.home_id, path: f.path, photo_key: f.photo_key, # ... 其余字段省略 }) |> Repo.all()
希望按home_id和path分组,返回每组按item_id排序的前3条结果,需要用Ecto实现类似SQL中row_number()分区的逻辑。
解决方案
可以通过Ecto的窗口函数row_number/0结合子查询来实现,具体步骤如下:
- 先合并两个子查询,并为每条记录添加分区后的行号
- 在外层查询中过滤出行号≤3的记录
完整代码实现
# 第一步:合并查询并添加行号 ranked_query = all_counts |> join(:inner, [c], f in ^file_info, on: c.home_id == f.home_id and c.path == f.path) |> select([c, f], %{ item_id: f.item_id, home_id: f.home_id, path: f.path, photo_key: f.photo_key, item_fruits_count: c.item_fruits_count, item_veggies_count: c.item_veggies_count, # 按home_id和path分区,按item_id排序生成行号 rn: row_number() |> over(partition_by: [f.home_id, f.path], order_by: [asc: f.item_id]) }) |> subquery() # 第二步:过滤前3条结果 result = ranked_query |> where([r], r.rn <= 3) |> select([r], %{ item_id: r.item_id, home_id: r.home_id, path: r.path, photo_key: r.photo_key, item_fruits_count: r.item_fruits_count, item_veggies_count: r.item_veggies_count }) |> Repo.all()
关键说明
- 使用
row_number() |> over(partition_by: [f.home_id, f.path], order_by: [asc: f.item_id])实现按两个字段分区,按item_id升序排序生成行号 - 将带行号的查询作为子查询,在外层过滤行号≤3的记录,确保每个
home_id+path分组只返回前3条数据 - 这种方式将所有计算逻辑推给数据库执行,避免了在内存中加载海量数据
内容的提问来源于stack exchange,提问作者Flame_Phoenix
相关产品推荐
相关产品推荐

