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

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)。

二、针对无报错崩溃的排查与修复

  1. 捕获子线程异常:当前子线程的异常仅在内部处理,主线程无法感知。需通过future.result()获取子线程异常,在主线程中记录:
    # 在threadpool.shutdown前添加
    try:
        futures.result()
    except Exception as e:
        logging.error(f'子线程写入错误: {e}', exc_info=True)
        print(f'子线程写入错误: {e}')
    
  2. 强制日志刷新:程序崩溃时日志可能滞留在缓冲区未写入文件,修改日志配置并在关键位置手动刷新:
    # 替换原日志配置
    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()
    
  3. 调整分批大小与连接超时:尝试调小ROW_BATCH_SIZE(比如从500改为200),同时在连接字符串中添加超时参数(如果驱动支持):
    cnxnString = r'Driver=FoundrySqlDriver;baseUrl='+server+';pwd='+tokenpwd+';Timeout=300;'
    
  4. 解决时序与队列积压问题:VSCode步进后能正常运行,说明主线程推送数据速度可能远超子线程处理速度,导致队列积压或连接中断。限制队列最大长度并设置阻塞等待:
    # 创建队列时设置最大长度
    data_queue = Queue(maxsize=10)
    # 推送数据时阻塞等待队列空闲
    data_queue.put(data_rows, block=True)
    
  5. 排查输出转换器问题:cnxn.add_output_converter(pyodbc.SQL_WVARCHAR, decode_sketchy_utf16)处理大文本时可能触发底层错误,暂时注释该代码,用默认解码方式测试。
  6. 添加全局异常捕获:在主入口捕获所有未处理异常,确保任何崩溃都能被记录:
    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)
    
  7. 检查驱动兼容性:尝试更新Palantir Foundry ODBC驱动到最新版本,查看驱动官方文档是否有Windows平台的特殊配置要求。

内容的提问来源于stack exchange,提问作者D W

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 01:02:54