Python通过pyodbc连接Palantir Foundry DB无报错突然终止排查求助
问题
使用pyodbc连接Windows平台上的Palantir Foundry自定义ODBC驱动,需导出约200万条含大文本字段的数据,因此采用fetchmany分批拉取而不使用fetchall。此前曾出现超时错误并已优化,但现在程序常拉取几千条数据后突然退出,无任何报错信息、超时提示,日志和输出也未写入结束标识。尝试过ODBC追踪(但日志写入导致速度过慢,仅能看到超时错误)、添加大量try-except块(包括包裹主函数调用)均无效。此外,在VSCode中步进前5轮左右后再让程序运行可正常完成,但直接从命令行或无断点运行时就会崩溃。
问题代码
import csv import datetime, time from string import ascii_uppercase import shutil, os, json, re import pyodbc import logging from concurrent.futures import ThreadPoolExecutor from queue import Queue #set vars server = r'<server>' tokenpwd = '<key>' cnxnString = r'Driver=FoundrySqlDriver;baseUrl='+server+';pwd='+tokenpwd+';' base_folder = 'C:\Python\' columns_file = 'clientData_export_columns.txt' count_file = 'results_count.json' output_file = f'C:\bin\clientData_Data.csv' QRY_STR = r'SELECT * FROM "<server table name>"' ROW_BATCH_SIZE = 500 log_file = base_folder + f'\LOGS\clientData_LOG_{time.strftime("%Y%m%d")}.log' today = datetime.datetime.now().date() yesterday = (today - datetime.timedelta(days=1)) tomorrow = (today + datetime.timedelta(days=1)) #today = today.strftime("%Y-%m-%d") #yesterday = yesterday.strftime("%Y-%m-%d") # Setup Logging logging.basicConfig(filename=log_file, filemode='a+', format='%(asctime)s - %(levelname)s - %(message)s', level=logging.DEBUG) def decode_sketchy_utf16(raw_bytes): s = raw_bytes.decode("utf-16-le", "ignore") # s = raw_bytes.decode("utf-8", "ignore") try: n = s.index('\u0000') s = s[:n] # respect null terminator except ValueError: pass return s def get_write_count(newcount = -1): with open(base_folder + count_file, "r+") as cfile: count_data = json.load(cfile) if newcount == -1: for c in count_data: if c['data'].upper() == 'clientData': return c['count'] else: for c in count_data: if c['data'].upper() == 'clientData': c['count'] = newcount cfile.seek(0) cfile.truncate() json.dump(count_data, cfile) break def write_data(data_queue, columns, columnList): try: logging.info('starting file write') with open(output_file, "w+", encoding="utf-8", newline="") as writecsv: writer = csv.DictWriter(writecsv, fieldnames=columnList, delimiter="|", quotechar='"', quoting=csv.QUOTE_ALL, extrasaction='ignore', escapechar='\\') writer.writeheader() while True: # print(f'queue size: {data_queue.qsize()}') data_rows = data_queue.get() if not data_rows or data_rows is None: break dataset = [dict(zip(columns, row)) for row in data_rows] writer.writerows(dataset) data_queue.task_done() except Exception as e: print(e) logging.error(f'Error writing file: {e}') logging.info('FILE WRITE ERROR(S)') print('Closed with error.') def export_clientData(): with open(base_folder + columns_file, "r") as col_file: columnList = col_file.readlines() columnList = [line.rstrip('\n') for line in columnList] get_write_count(0) print('Exporting...') logging.info('BEGIN EXPORT') try: cnxn = pyodbc.connect(cnxnString) with cnxn: cursor = cnxn.cursor() try: cursor.execute(QRY_STR) x = 0 # logging.info('Executed qry.') except pyodbc.Error as ex: sqlstate = ex.args[1] # print(sqlstate) logging.error(f'qry: {QRY_STR}') logging.error(f'sql_err: {sqlstate}') logging.error('Error executing query') raise Exception # cursor.execute("ROLLBACK") except Exception as e: print(e) logging.error(f'Some cursor execution error: {e}') logging.info('END EXPORT WITH ERROR(S)') print('Closed with error.') return 1 columns = [column[0] for column in cursor.description] TestExportDate = True # boolean for whether we need to test Exported Date rcounter = 0 cnxn.add_output_converter(pyodbc.SQL_WVARCHAR, decode_sketchy_utf16) data_rows = cursor.fetchmany(ROW_BATCH_SIZE) data_queue = Queue() dataset = [dict(zip(columns, row)) for row in data_rows] if TestExecutedDate == True: TestExecutedDate = False ex_date = dataset[0]['CLIENT_EXPORT_DATE'] print(f'Data CLIENT_EXPORT_DATE date: {ex_date}') if dataset[0]['CLIENT_EXPORT_DATE'] not in [yesterday, today, tomorrow]: os.remove(output_file) print(f'CLIENT_EXPORT_DATE did not fall on current or previous dates ({yesterday} or {today})') logging.error(f'CLIENT_EXPORT_DATE did not fall on current or previous dates ({yesterday} or {today})') return 1 with ThreadPoolExecutor(max_workers=1) as threadpool: futures = threadpool.submit(write_data, data_queue, columns, columnList) while data_rows: data_queue.put(data_rows) data_rows = cursor.fetchmany(ROW_BATCH_SIZE) if not data_rows: # or x >= 25: data_queue.put(None) break rcounter += len(data_rows) x += 1 if x % 250000 == 0: print(f'rows: {rcounter}') get_write_count(rcounter) threadpool.shutdown(wait=True, cancel_futures=False) # future.result() # data_queue.put(None) except Exception as e: print(e) logging.error(f'Error creating odbc connection: {e}') logging.info('END EXPORT WITH ERROR(S)') print('Closed with error.') return 1 get_write_count(rcounter) get_write_count(executed=f'{ex_date}') logging.info('END DATA EXPORT') print('Finished data Export.') return 0 if __name__ == '__main__': export_clientData()
解决方案
一、先修复代码中的明显错误
- 修正变量名错误:代码中定义了
TestExportDate = True,但后续判断用了未定义的TestExecutedDate,统一改为TestExportDate,避免触发未捕获的NameError。 - 修复路径转义问题:所有Windows路径改用原始字符串(前缀加
r),比如base_folder = r'C:\Python'、output_file = r'C:\bin\clientData_Data.csv',避免反斜杠转义导致的路径错误或语法错误。 - 修复无效函数调用:
get_write_count(executed=f'{ex_date}')是无效调用,该函数仅接受newcount参数,需删除此调用或修改函数定义支持日期记录。 - 修正行数统计逻辑:第一批拉取的
data_rows长度未计入rcounter,需在循环开始前添加rcounter += len(data_rows)。
二、针对无报错崩溃的排查与修复
- 捕获子线程异常:当前子线程的异常仅在内部处理,主线程无法感知。需通过
future.result()获取子线程异常,在主线程中记录:# 在threadpool.shutdown前添加 try: futures.result() except Exception as e: logging.error(f'子线程写入错误: {e}', exc_info=True) print(f'子线程写入错误: {e}') - 强制日志刷新:程序崩溃时日志可能滞留在缓冲区未写入文件,修改日志配置并在关键位置手动刷新:
# 替换原日志配置 handler = logging.FileHandler(log_file, mode='a+', encoding='utf-8') formatter = logging.Formatter('%(asctime)s - %(levelname)s - %(message)s') handler.setFormatter(formatter) logging.basicConfig(level=logging.DEBUG, handlers=[handler], force=True) # 在异常处理、循环计数后添加 logging.flush() - 调整分批大小与连接超时:尝试调小
ROW_BATCH_SIZE(比如从500改为200),同时在连接字符串中添加超时参数(如果驱动支持):cnxnString = r'Driver=FoundrySqlDriver;baseUrl='+server+';pwd='+tokenpwd+';Timeout=300;' - 解决时序与队列积压问题:VSCode步进后能正常运行,说明主线程推送数据速度可能远超子线程处理速度,导致队列积压或连接中断。限制队列最大长度并设置阻塞等待:
# 创建队列时设置最大长度 data_queue = Queue(maxsize=10) # 推送数据时阻塞等待队列空闲 data_queue.put(data_rows, block=True) - 排查输出转换器问题:
cnxn.add_output_converter(pyodbc.SQL_WVARCHAR, decode_sketchy_utf16)处理大文本时可能触发底层错误,暂时注释该代码,用默认解码方式测试。 - 添加全局异常捕获:在主入口捕获所有未处理异常,确保任何崩溃都能被记录:
if __name__ == '__main__': try: result = export_clientData() except Exception as e: logging.error(f'全局未捕获异常: {e}', exc_info=True) print(f'全局未捕获异常: {e}') exit(1) exit(result) - 检查驱动兼容性:尝试更新Palantir Foundry ODBC驱动到最新版本,查看驱动官方文档是否有Windows平台的特殊配置要求。
内容的提问来源于stack exchange,提问作者D W
相关产品推荐
相关产品推荐

