Python使用Spanner带参数执行分页查询时陷入停滞无报错
问题
在Python中使用Spanner实现分页查询,通过LIMIT和OFFSET结合execute_sql的params参数传参做迭代分页时,程序陷入停滞,无任何错误提示。具体情况:
- 将OFFSET和LIMIT的值直接硬编码进查询语句,执行成功;
- 已显式指定
params_types确保参数类型正确,且其他查询用params传参正常; - 将查询中的
INNER JOIN改为LEFT JOIN后,代码可正常运行; - 该查询在Spanner控制台能成功执行并返回结果。
测试代码如下:
def print_test(namespace, field_id): try: table_id = f"{namespace.lower().strip()}_{field_id.lower().strip()}" page_size = 50000 offset_param = 0 while True: with database.snapshot() as snapshot: params = { "offset_param": offset_param, "page_size": page_size } query = f""" SELECT P.ML_ID, S.VALUE, S.SITE, S.EXECUTION_ID FROM `sync_{table_id}` S INNER JOIN `people` P ON P.CUST_ID = CAST(S.VALUE AS INT64) WHERE S.ML_ID IS NULL LIMIT @page_size OFFSET @offset_param; """ params_types = { "offset_param": spanner.param_types.INT64, "page_size": spanner.param_types.INT64 } print(params) print(query) results = snapshot.execute_sql(query, params=params, param_types = params_types ) records_to_process = list(results) if not records_to_process: break print(records_to_process) offset_param += page_size except Exception as e: print(f"Error occurred: {e}") print_test(namespace, field_id)
分析与解决方案
1. 核心原因
参数化的OFFSET/LIMIT结合INNER JOIN时,Spanner的查询优化器可能生成了低效的执行计划——需要扫描大量无关数据后再做JOIN和分页,导致查询长时间运行(表现为程序停滞)。而硬编码参数或改用LEFT JOIN时,优化器能生成更优的执行路径,因此执行正常。
2. 优化方案
方案一:强制使用索引引导查询计划
检查sync_{table_id}和people表是否有适配的索引(比如针对S.ML_ID、S.VALUE或P.CUST_ID的索引),然后在查询中使用FORCE_INDEX提示,强制优化器选择高效路径:
SELECT P.ML_ID, S.VALUE, S.SITE, S.EXECUTION_ID FROM `sync_{table_id}` S FORCE_INDEX(idx_sync_ml_id_value) INNER JOIN `people` P FORCE_INDEX(idx_people_cust_id) ON P.CUST_ID = CAST(S.VALUE AS INT64) WHERE S.ML_ID IS NULL LIMIT @page_size OFFSET @offset_param;
方案二:改用键集分页替代OFFSET分页
Spanner官方推荐使用键集分页而非OFFSET,因为OFFSET在数据量大时会因扫描并跳过前置行导致性能暴跌。实现步骤:
- 选择唯一有序列(如
S.EXECUTION_ID)作为分页标记; - 以上一页最后一条记录的键值作为下一页的起始条件,替代
OFFSET。
修改后的示例代码:
def print_test(namespace, field_id): try: table_id = f"{namespace.lower().strip()}_{field_id.lower().strip()}" page_size = 50000 last_execution_id = None while True: with database.snapshot() as snapshot: params = {"page_size": page_size} query_base = """ SELECT P.ML_ID, S.VALUE, S.SITE, S.EXECUTION_ID FROM `sync_{table_id}` S INNER JOIN `people` P ON P.CUST_ID = CAST(S.VALUE AS INT64) WHERE S.ML_ID IS NULL """ # 添加分页条件 if last_execution_id is not None: query = query_base + " AND S.EXECUTION_ID > @last_execution_id " params["last_execution_id"] = last_execution_id params_types = { "last_execution_id": spanner.param_types.STRING, # 根据实际字段类型调整 "page_size": spanner.param_types.INT64 } else: query = query_base params_types = {"page_size": spanner.param_types.INT64} query += " ORDER BY S.EXECUTION_ID LIMIT @page_size;" results = snapshot.execute_sql(query, params=params, param_types=params_types) records_to_process = list(results) if not records_to_process: break print(records_to_process) # 更新分页标记 last_execution_id = records_to_process[-1][3] except Exception as e: print(f"Error occurred: {e}")
方案三:拆分查询逻辑减少JOIN数据量
先从sync_{table_id}筛选符合条件的记录,再批量与people表关联,降低单次JOIN的数据规模:
def print_test(namespace, field_id): try: table_id = f"{namespace.lower().strip()}_{field_id.lower().strip()}" page_size = 50000 offset_param = 0 while True: with database.snapshot() as snapshot: # 第一步:批量获取sync表符合条件的记录 sync_params = {"offset_param": offset_param, "page_size": page_size} sync_query = """ SELECT VALUE, SITE, EXECUTION_ID FROM `sync_{table_id}` WHERE ML_ID IS NULL LIMIT @page_size OFFSET @offset_param; """ sync_results = snapshot.execute_sql(sync_query, params=sync_params, param_types={ "offset_param": spanner.param_types.INT64, "page_size": spanner.param_types.INT64 }) sync_records = list(sync_results) if not sync_records: break # 第二步:批量查询people表匹配数据 value_list = [int(record[0]) for record in sync_records] people_query = """ SELECT CUST_ID, ML_ID FROM `people` WHERE CUST_ID IN UNNEST(@value_list); """ people_results = snapshot.execute_sql(people_query, params={"value_list": value_list}, param_types={ "value_list": spanner.param_types.Array(spanner.param_types.INT64) }) people_map = {row[0]: row[1] for row in people_results} # 第三步:关联结果 final_records = [] for sync_record in sync_records: cust_id = int(sync_record[0]) if cust_id in people_map: final_records.append((people_map[cust_id], sync_record[0], sync_record[1], sync_record[2])) print(final_records) offset_param += page_size except Exception as e: print(f"Error occurred: {e}")
3. 临时调试方案
开启Spanner查询日志,对比参数化与非参数化查询的执行计划差异:
- 在代码中启用Spanner客户端调试日志;
- 在Spanner控制台查看该查询的执行详情,分析耗时瓶颈。
内容的提问来源于stack exchange,提问作者Miller Lizarazo
相关产品推荐
相关产品推荐

