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

多线程数据爬取实现疑问:运行出现‘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)

分析与解决

多线程实现的核心问题

你的多线程实现存在几个关键错误,同时这些问题也是导致进程被杀死的诱因:

  1. executor.map传参错误:links_to_process是元组列表,但executor.map会把每个元组当成单个参数传给fetch_url_and_process,导致函数接收到的是(link_url, serial_number)作为第一个参数,第二个参数缺失,触发未捕获的参数异常,直接导致线程崩溃。
  2. 日志未配置:代码仅导入logging但未设置输出配置,即使有错误也不会打印到控制台,导致你看不到崩溃前的异常信息。
  3. 异常捕获范围过窄:仅捕获Oracle相关错误,爬取、文件操作、S3上传等环节的异常会直接终止线程且无日志记录,隐藏问题。
  4. 批量加载数据导致内存溢出:如果Link_Status=0的条目数量极大,fetchall()会一次性把所有数据加载到内存,导致内存占用飙升,这是系统触发OOM(内存不足)杀死进程的核心原因(控制台打印Killed是Linux系统OOM Killer的典型行为)。

修复建议

  1. 修正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)}")
    
  2. 配置日志输出
    在代码开头添加日志配置,让异常信息能显示:

    logging.basicConfig(
        level=logging.ERROR,
        format='%(asctime)s - %(levelname)s - %(message)s',
        handlers=[logging.StreamHandler()]
    )
    
  3. 扩大异常捕获范围
    在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)
    
  4. 分批获取待爬取链接
    避免一次性加载所有数据到内存,改用分页查询:

    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
    
  5. 优化内存占用
    爬取大页面时限制文本大小,避免内存堆积:

    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 21:07:33