多线程数据爬取实现疑问:运行出现‘Killed’终止故障
问题
我尝试为数据爬取实现多线程功能,使用的Oracle表DATA_TABLE包含serial_number(追踪条目数)、Link_URL(待爬取链接)、Link_Status(0=未爬取,1=已爬取)三列。运行下方代码后,控制台在10-15分钟后仅打印Killed就终止,无其他日志输出,想确认当前多线程实现是否正确。
代码:
import os import boto3 import cx_Oracle import logging import configparser from trafilatura import fetch_url, extract from concurrent.futures import ThreadPoolExecutor # Load configuration from the config.ini file config = configparser.ConfigParser() config.read('config.ini') oracle_username = config.get('database', 'username') oracle_password = config.get('database', 'password') oracle_host = config.get('database', 'host') oracle_port = config.get('database', 'port') oracle_service_name = config.get('database', 'service_name') # Create a connection pool to ensure each thread has its own separate connection pool = cx_Oracle.SessionPool( user=oracle_username, password=oracle_password, dsn=f"{oracle_host}:{oracle_port}/{oracle_service_name}", min=2, max=4, increment=1, threaded=True ) s3_client = boto3.client('s3') bucket_name = config.get('S3', 's3_bucket_name') folder_path = config.get('S3', 's3_folder_path_test') def fetch_url_and_process(link_url, serial_number): try: connection = pool.acquire() with connection.cursor() as cursor: # Scrap the text from the URL using trafilature downloaded = fetch_url(link_url) extracted_text = extract(downloaded) print('******************************************') print(extracted_text) print('******************************************') if extracted_text is None: extracted_text = 'delete' # Save the extracted text to a text file local_file_name = f'{serial_number}.txt' print('--------------------------->', local_file_name) with open(local_file_name, 'w', encoding='utf-8') as file: file.write(extracted_text) # Upload the text file to S3 bucket s3_client.upload_file(local_file_name, bucket_name, f'{folder_path}{local_file_name}') print("Text file uploaded to S3") # Delete the local text file os.remove(local_file_name) # Update Link_Status to 1 update_query = "UPDATE DATA_TABLE SET Link_Status = 1 WHERE Link_URL = :url" cursor.execute(update_query, {'url': link_url}) connection.commit() except cx_Oracle.Error as e: logging.error(f"Error while inserting data into Oracle: {e}") finally: pool.release(connection) if __name__ == "__main__": try: connection = pool.acquire() with connection.cursor() as cursor: # Fetch all Link_URLs where Link_Status = 0 query = "SELECT Link_URL, Serial_Number FROM DATA_TABLE WHERE Link_Status = 0" cursor.execute(query) rows = cursor.fetchall() if rows: links_to_process = [(row[0], row[1]) for row in rows] with ThreadPoolExecutor(max_workers=4) as executor: executor.map(fetch_url_and_process, links_to_process) except cx_Oracle.Error as e: logging.error(f"Error while fetching data from Oracle: {e}") finally: pool.release(connection)
分析与解决
多线程实现的核心问题
你的多线程实现存在几个关键错误,同时这些问题也是导致进程被杀死的诱因:
executor.map传参错误:links_to_process是元组列表,但executor.map会把每个元组当成单个参数传给fetch_url_and_process,导致函数接收到的是(link_url, serial_number)作为第一个参数,第二个参数缺失,触发未捕获的参数异常,直接导致线程崩溃。- 日志未配置:代码仅导入
logging但未设置输出配置,即使有错误也不会打印到控制台,导致你看不到崩溃前的异常信息。 - 异常捕获范围过窄:仅捕获Oracle相关错误,爬取、文件操作、S3上传等环节的异常会直接终止线程且无日志记录,隐藏问题。
- 批量加载数据导致内存溢出:如果
Link_Status=0的条目数量极大,fetchall()会一次性把所有数据加载到内存,导致内存占用飙升,这是系统触发OOM(内存不足)杀死进程的核心原因(控制台打印Killed是Linux系统OOM Killer的典型行为)。
修复建议
修正
executor.map传参方式
改用拆分参数的方式传递,或者用executor.submit:# 方式1:拆分元组传递 urls, serial_numbers = zip(*links_to_process) with ThreadPoolExecutor(max_workers=4) as executor: executor.map(fetch_url_and_process, urls, serial_numbers) # 方式2:用submit并捕获任务异常 with ThreadPoolExecutor(max_workers=4) as executor: futures = [executor.submit(fetch_url_and_process, url, sn) for url, sn in links_to_process] for future in concurrent.futures.as_completed(futures): try: future.result() except Exception as e: logging.error(f"Task failed: {str(e)}")配置日志输出
在代码开头添加日志配置,让异常信息能显示:logging.basicConfig( level=logging.ERROR, format='%(asctime)s - %(levelname)s - %(message)s', handlers=[logging.StreamHandler()] )扩大异常捕获范围
在fetch_url_and_process中捕获所有类型的异常,避免线程静默崩溃:try: # 原有业务代码 except cx_Oracle.Error as e: logging.error(f"Oracle error processing {link_url}: {str(e)}") except Exception as e: logging.error(f"Unexpected error processing {link_url}: {str(e)}") finally: pool.release(connection)分批获取待爬取链接
避免一次性加载所有数据到内存,改用分页查询:batch_size = 100 offset = 0 while True: query = """ SELECT Link_URL, Serial_Number FROM DATA_TABLE WHERE Link_Status = 0 FETCH FIRST :batch ROWS ONLY OFFSET :offset ROWS """ cursor.execute(query, {'batch': batch_size, 'offset': offset}) rows = cursor.fetchall() if not rows: break urls, serial_numbers = zip(*rows) with ThreadPoolExecutor(max_workers=4) as executor: executor.map(fetch_url_and_process, urls, serial_numbers) offset += batch_size优化内存占用
爬取大页面时限制文本大小,避免内存堆积:extracted_text = extract(downloaded) if extracted_text is None: extracted_text = 'delete' else: # 截断过长文本,根据业务调整阈值 extracted_text = extracted_text[:100000]
关于Killed的说明
控制台打印Killed几乎都是系统OOM Killer的行为,根源是进程占用内存超过系统阈值。按上述建议分批处理数据、优化内存后,该问题会得到解决。
内容的提问来源于stack exchange,提问作者Vinay Sharma
相关产品推荐
相关产品推荐

