Python 다중 스레드를 활용한 아마존 상품 정보 추출

개발 환경 설정

아마존 상품 데이터를 효율적으로 수집하기 위해 독립적인 개발 환경을 구축하는 것이 중요합니다. 여기서는 miniconda를 사용하여 새로운 파이썬 환경을 생성하고 필요한 라이브러리를 설치하는 과정을 설명합니다. 이는 다른 프로젝트와의 의존성 충돌을 방지하는 데 도움이 됩니다.

Conda 환경 생성 및 활성화

# 'amazon_scraper_env'라는 이름으로 Python 3.12 환경 생성
conda create -n amazon_scraper_env python=3.12

# 생성된 환경 목록 확인
conda info --envs

# 새 환경 활성화
conda activate amazon_scraper_env

필수 라이브러리 설치

환경이 활성화되면, 데이터 수집에 필요한 패키지들을 pip를 사용하여 설치합니다. 포함된 라이브러리는 pymysql (데이터베이스 연동), requests (HTTP 요청), retrying (재시도 로직), lxml (HTML 파싱), loguru (로깅), feapder (크롤링 프레임워크) 등입니다.

pip install pymysql requests retrying lxml loguru feapder

환경 관리 명령어

설치된 패키지 확인, 환경 비활성화 및 삭제를 위한 명령어는 다음과 같습니다.

# pip로 설치된 패키지 목록 확인
pip list

# 현재 환경 비활성화
conda deactivate

# 'amazon_scraper_env' 환경 및 모든 패키지 삭제 (선택 사항)
conda remove -n amazon_scraper_env --all

IDE(예: PyCharm)에서 새롭게 생성된 Conda 환경의 Python 인터프리터를 프로젝트에 맞게 설정하여 사용할 수 있습니다.

아마존 카테고리 정보 수집

아마존 웹사이트(https://www.amazon.com)의 상품 카테고리는 일반적으로 동적으로 로드됩니다. 개발자 도구를 사용하여 네트워크 트래픽을 분석하면, 메인 내비게이션 메뉴를 위한 특정 AJAX 요청 URL을 식별할 수 있습니다.

예시 URL: https://www.amazon.com/nav/ajax/hamburgerMainContent?... 이와 같은 엔드포인트는 일반적으로 HTML 형식으로 모든 메뉴 계층 구조를 반환합니다.

HTML 콘텐츠 파싱 및 카테고리 추출

반환된 HTML에서 상품 카테고리 링크를 추출하기 위해 XPath를 활용합니다. 특정 data-menu-id 범위에 있는 ul 요소 내에서, 첫 두 개의 li 요소를 제외하고 hmenu-separator 클래스를 가지지 않는 li 요소를 대상으로 링크를 추출할 수 있습니다.

//ul[number(@data-menu-id) >= 5 and number(@data-menu-id) <= 26]/li[position() > 2 and not(contains(@class, "hmenu-separator"))]/a/@href

이 XPath는 상품 카테고리 섹션에 해당하는 링크(href 속성)만을 효과적으로 찾아냅니다. 추출된 링크는 상대 경로일 수 있으므로, 완전한 URL을 만들기 위해 아마존 도메인을 앞에 추가해야 합니다.

다국어 콘텐츠 요청 처리

아마존은 사용자의 지역 설정에 따라 다른 언어의 콘텐츠를 제공합니다. 한국어 또는 특정 언어의 콘텐츠를 얻기 위해서는 HTTP 요청 헤더에 Cookie 및 Accept-Language 정보를 포함시켜야 합니다.

headers = {
    "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36",
    "Cookie": "i18n-prefs=KRW; lc-main=ko_KR;", # 통화 및 언어 설정
    "Accept-Language": "ko-KR,ko;q=0.9,en-US;q=0.8,en;q=0.7",
}
# requests.get(url, headers=headers) 호출 시 이 헤더를 사용합니다.

초기 카테고리 링크 수집은 일반적으로 단일 스레드로 충분하며, 이후 단계에서 다중 스레딩을 적용하여 효율성을 높일 수 있습니다.

상품 목록 및 상세 페이지 링크 추출

카테고리 링크를 통해 접근한 상품 목록 페이지에서는 일반적으로 여러 상품과 함께 페이지네이션(pagination)이 존재합니다. 우리는 페이지네이션 정보를 분석하여 전체 상품 목록을 탐색하고, 각 상품의 상세 페이지 링크를 추출해야 합니다.

페이지네이션 처리

상품 목록 페이지에서 가장 큰 페이지 번호를 찾아, 이를 기반으로 모든 페이지의 URL을 생성합니다. 일반적으로 URL 쿼리 파라미터(예: &page=)를 조작하여 각 페이지에 접근할 수 있습니다. 생성된 모든 페이지 URL은 상품 상세 링크 추출 작업을 담당하는 스레드에 전달됩니다.

상품 상세 링크 추출

각 상품 목록 페이지에서는 개별 상품으로 연결되는 링크를 추출합니다. 이 링크들은 일반적으로 상품 이미지나 제목을 감싸는 <a> 태그의 href 속성에 포함되어 있습니다.

//div[@data-component-type="s-search-result"]//h2/a/@href

위 XPath 예시는 아마존 검색 결과 페이지에서 상품 제목 링크를 추출하는 데 사용될 수 있습니다. 추출된 상품 상세 URL은 별도의 큐(product_url_queue)에 저장되어 상품 정보 파싱 작업을 기다립니다.

큐에서 작업을 가져올 때는 get() 메서드를 사용하고, 작업 완료 후에는 반드시 task_done() 메서드를 호출하여 큐에 있는 총 작업 수를 줄여야 합니다. 이는 모든 작업이 완료될 때까지 기다리는 join() 메서드의 정확한 작동을 보장합니다.

상품 상세 정보 파싱

상품 상세 페이지는 제품의 핵심 정보를 담고 있습니다. 여기서는 상품 제목, 이미지 URL, 가격, 판매량, 고객 리뷰 수 등 필요한 데이터를 추출하는 과정을 다룹니다. 각 상품 유형에 따라 페이지 구조가 다소 상이할 수 있으므로, 견고한 XPath 또는 CSS 선택자를 사용하는 것이 중요합니다.

주요 데이터 추출

상품 제목과 같은 핵심 정보는 고유한 ID나 클래스를 가진 요소에 위치하는 경우가 많습니다. 예를 들어, 상품 제목은 다음과 같은 XPath로 추출할 수 있습니다.

//span[@id="productTitle"]/text()

데이터 추출 시에는 XPath 결과가 없을 경우를 대비하여 항상 빈 값 확인 로직을 포함하여 코드의 견고성을 높여야 합니다.

import requests
from lxml import html
import time
from retrying import retry

class ProductScraper:
    def __init__(self):
        self.headers = {
            "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36",
            "Accept-Language": "ko-KR,ko;q=0.9,en-US;q=0.8,en;q=0.7",
        }

    @retry(stop_max_attempt_number=3, wait_fixed=2000)
    def fetch_page_content(self, url):
        print(f"Fetching: {url}")
        response = requests.get(url, headers=self.headers, timeout=10)
        response.raise_for_status() # HTTP 오류 발생 시 예외 발생
        tree = html.fromstring(response.content)
        # 예시: 상품 제목이 존재하는지 확인하여 유효성 검사
        if not tree.xpath('//span[@id="productTitle"]'):
            raise Exception("Product title not found, retrying...")
        return tree

    def parse_product_data(self, product_url):
        try:
            tree = self.fetch_page_content(product_url)
            product_title = self.extract_text(tree, '//span[@id="productTitle"]/text()')
            product_price = self.extract_text(tree, '//span[contains(@class, "a-price")]/span[@aria-hidden="true"]/text()')
            # 다른 데이터 추출 로직...
            return {
                "title": product_title.strip() if product_title else "N/A",
                "price": product_price.strip() if product_price else "N/A",
                "url": product_url
            }
        except Exception as e:
            print(f"Error parsing {product_url}: {e}")
            return None

    def extract_text(self, tree, xpath_expression):
        elements = tree.xpath(xpath_expression)
        return elements[0] if elements else None

# 사용 예시 (주석 처리됨)
# scraper = ProductScraper()
# data = scraper.parse_product_data("https://www.amazon.com/dp/B0BP7M512P")
# if data:
#     print(data)

위 코드 예시에서는 retrying 라이브러리를 사용하여 페이지 요청 및 특정 요소 유무 검사 시 재시도 로직을 구현했습니다. 이는 네트워크 문제나 일시적인 안티-크롤링 조치로 인해 데이터 추출이 실패할 경우 유용합니다. 재시도는 최대 3번까지 수행하며, 각 시도 사이에 2초의 지연 시간을 둡니다.

수집 데이터 저장 (MySQL)

수집된 상품 데이터는 MySQL 데이터베이스에 저장될 수 있습니다. 데이터베이스에 데이터를 삽입할 때, 개별 레코드를 하나씩 삽입하는 것보다 여러 레코드를 한 번에 삽입하는 배치 삽입(Batch Insertion) 방식을 사용하는 것이 성능 면에서 훨씬 효율적입니다.

아래는 pymysql을 사용하여 데이터를 배치로 삽입하는 예시 코드입니다. 여기서는 30개의 레코드를 모아 한 번에 삽입하도록 구성되어 있습니다. 이 방식은 데이터베이스 부하를 줄이지만, 삽입 실패 시 배치 단위로 데이터가 손실될 위험도 고려해야 합니다.

import pymysql
import pymysql.cursors
from datetime import datetime

class ProductDBManager:
    def __init__(self, host, user, password, db_name):
        self.conn_params = {
            'host': host,
            'user': user,
            'password': password,
            'database': db_name,
            'charset': 'utf8mb4',
            'cursorclass': pymysql.cursors.DictCursor # 딕셔너리 형태로 결과 반환
        }
        self.products_to_save = []
        self.batch_size = 30

    def add_product_data(self, product_data):
        self.products_to_save.append(product_data)
        if len(self.products_to_save) >= self.batch_size:
            self.flush_data()

    def flush_data(self):
        if not self.products_to_save:
            return

        sql_insert = """
            INSERT INTO products (title, price, url, extracted_at)
            VALUES (%s, %s, %s, %s)
        """
        try:
            with pymysql.connect(**self.conn_params) as conn:
                with conn.cursor() as cursor:
                    records = []
                    for p in self.products_to_save:
                        # 데이터에 현재 시간 추가
                        records.append((p['title'], p['price'], p['url'], datetime.now()))
                    cursor.executemany(sql_insert, records)
                conn.commit()
            print(f"{len(self.products_to_save)} records successfully inserted.")
            self.products_to_save.clear() # 성공적으로 삽입 후 목록 비우기
        except pymysql.Error as e:
            print(f"Error during batch insertion: {e}")
            # 오류 발생 시 데이터 손실 방지를 위한 추가 로직 필요 (예: 로깅, 재시도 큐)

    def close(self):
        # 애플리케이션 종료 시 잔여 데이터 저장
        self.flush_data()

# 사용 예시 (주석 처리됨)
# db_manager = ProductDBManager('localhost', 'root', 'password', 'amazon_data')
# db_manager.add_product_data({"title": "Test Product 1", "price": "$10.00", "url": "http://example.com/p1"})
# db_manager.add_product_data({"title": "Test Product 2", "price": "$20.00", "url": "http://example.com/p2"})
# # ... 더 많은 데이터 추가
# db_manager.close() # 종료 시 남은 데이터 저장

SQL 쿼리 작성 시 들여쓰기나 구문 오류는 프로그램 작동에 문제를 일으킬 수 있으므로 주의 깊게 확인해야 합니다.

다중 스레드를 활용한 작업 분배

대량의 데이터를 효율적으로 수집하기 위해 파이썬의 threading 모듈과 queue 모듈을 활용하여 다중 스레드 기반의 크롤링 시스템을 구축할 수 있습니다. 각기 다른 역할을 수행하는 스레드들을 생성하고, 큐를 통해 작업과 데이터를 주고받는 방식입니다.

스레드 역할 분배

작업의 특성과 중요도를 고려하여 각 단계에 필요한 스레드 수를 할당합니다.

  • 카테고리 링크 수집: 1개의 스레드 (초기 작업이므로 단일 스레드로 충분)
  • 상품 목록 페이지 처리 및 상세 링크 추출: 10개의 스레드 (병렬 처리로 효율 증대)
  • 상품 상세 정보 파싱: 5개의 스레드 (데이터 추출 및 정제 작업)
  • 데이터 저장: 1개의 스레드 (데이터베이스 트랜잭션 관리)

스레드 관리 및 실행

모든 스레드는 메인 스크립트에서 생성되고 시작됩니다. threading.Thread 객체를 사용하여 각 스레드를 정의하고, start() 메서드로 실행합니다. 이때 daemon=True로 설정하여 메인 스레드가 종료될 때 모든 자식 스레드가 자동으로 종료되도록 하는 것이 일반적입니다.

import threading
import queue
import time

# 가정: ProductScraper, ProductDBManager 클래스는 이전에 정의되어 있다고 가정합니다.

# 각 작업 큐 생성
category_link_queue = queue.Queue() # 초기 카테고리 URL
listing_page_queue = queue.Queue()  # 상품 목록 페이지 URL
detail_page_queue = queue.Queue()   # 상품 상세 페이지 URL
parsed_data_queue = queue.Queue()   # 파싱된 상품 데이터

# worker_task 함수: 큐에서 작업을 가져와 처리하고, 완료되면 task_done 호출
def worker_task(input_queue, processor_func, output_queue=None, *args):
    while True:
        try:
            task = input_queue.get(timeout=1) # 1초 동안 기다림
            result = processor_func(task, *args) # 실제 작업 처리
            if output_queue and result is not None:
                output_queue.put(result) # 다음 큐에 결과 전달
            input_queue.task_done()
        except queue.Empty:
            break # 큐가 비어있으면 종료
        except Exception as e:
            print(f"Worker error processing task {task}: {e}")
            input_queue.task_done() # 오류 발생해도 작업 완료 처리

# 플레이스홀더 함수 (실제 스크래퍼 로직으로 대체 필요)
def fetch_initial_categories(queue_to_fill):
    print("Fetching initial category URLs...")
    # 실제 아마존 햄버거 메뉴 API 호출 및 파싱 로직
    # 예시:
    queue_to_fill.put("https://www.amazon.com/s?k=electronics")
    queue_to_fill.put("https://www.amazon.com/s?k=books")
    print("Initial categories added to queue.")

def process_listing_page(url, queue_to_fill):
    print(f"Processing listing page: {url}")
    # 실제 requests.get(url) 및 lxml 파싱 로직
    # 페이지네이션 처리 및 상세 URL 추출
    # 예시:
    for i in range(1, 3): # 두 개의 가상 상세 페이지 링크 생성
        detail_url = f"{url.replace('s?k=', 'dp/sample_')}_p{i}"
        queue_to_fill.put(detail_url)
    return None # 이 단계에서는 직접적인 결과를 반환하지 않음

def extract_product_details(detail_url):
    print(f"Extracting details for: {detail_url}")
    # ProductScraper 인스턴스를 사용하여 상세 정보 파싱
    # 예시 (실제 ProductScraper 클래스 인스턴스 사용):
    # scraper = ProductScraper()
    # data = scraper.parse_product_data(detail_url)
    # return data
    
    # 가상 데이터 반환
    return {"title": f"Product Title for {detail_url.split('/')[-1]}", "price": "$XX.XX", "url": detail_url}

def save_processed_data(data):
    print(f"Saving data: {data['title']}")
    # ProductDBManager 인스턴스를 사용하여 데이터 저장
    # 예시 (실제 ProductDBManager 클래스 인스턴스 사용):
    # db_manager = ProductDBManager('localhost', 'root', 'password', 'amazon_data')
    # db_manager.add_product_data(data)
    return None # 데이터 저장 작업은 다음 큐로 전달할 결과가 없음


# 스레드 리스트
worker_threads = []

# 1. 초기 카테고리 수집 (단일 스레드)
initial_fetch_thread = threading.Thread(target=fetch_initial_categories, args=(listing_page_queue,), daemon=True)
worker_threads.append(initial_fetch_thread)

# 2. 상품 목록 페이지 처리 스레드 (10개)
for _ in range(10):
    t = threading.Thread(target=worker_task, args=(listing_page_queue, process_listing_page, detail_page_queue), daemon=True)
    worker_threads.append(t)

# 3. 상품 상세 정보 파싱 스레드 (5개)
for _ in range(5):
    t = threading.Thread(target=worker_task, args=(detail_page_queue, extract_product_details, parsed_data_queue), daemon=True)
    worker_threads.append(t)

# 4. 데이터 저장 스레드 (1개)
# 이 스레드는 데이터를 DB에 저장하고, 다음 큐로 전달할 것이 없으므로 output_queue는 None
t_save_data = threading.Thread(target=worker_task, args=(parsed_data_queue, save_processed_data, None), daemon=True)
worker_threads.append(t_save_data)

# 모든 스레드 시작
for t in worker_threads:
    t.start()

# 초기 카테고리 수집 스레드를 시작하여 첫 작업 큐에 항목을 넣음
initial_fetch_thread.start()


# 모든 큐의 작업이 완료될 때까지 메인 스레드 대기
# 큐에 초기 작업이 모두 추가된 후 join 호출
# 각 큐의 모든 item들이 task_done()될 때까지 기다립니다.
listing_page_queue.join()
detail_page_queue.join()
parsed_data_queue.join()

print("All scraping tasks completed.")

모든 스레드가 실행되기 시작한 후, queue.join() 메서드를 호출하여 메인 스레드를 블록하고, 모든 큐에 들어있는 작업들이 완료될 때까지 기다립니다. 이는 크롤링 작업이 완전히 종료되는 것을 보장합니다.

프록시 IP 풀 활용 (선택 사항)

동일한 IP 주소로 아마존과 같은 웹사이트에 장시간, 고빈도로 요청을 보내면 IP 차단(block) 또는 접근 제한을 당할 수 있습니다. 이를 회피하고 안정적인 데이터 수집을 유지하기 위해 프록시 IP 풀을 사용하는 것이 좋습니다.

프록시 IP 풀은 유료 서비스로 제공되는 경우가 많으며, 일반적으로 다음과 같은 방식으로 시스템에 통합할 수 있습니다:

  1. 프록시 큐 초기화: 스크래퍼 초기화 시 프록시 서버 URL과 프록시 IP를 저장할 큐(proxy_queue)를 선언합니다.
  2. 프록시 IP 요청 및 관리 스레드: 별도의 스레드를 생성하여 프록시 서비스 제공자로부터 주기적으로 새로운 IP 주소를 가져와 proxy_queue에 추가합니다. 이 스레드는 사용 가능한 IP가 충분히 유지되도록 관리합니다.
  3. 스크래핑 요청에 프록시 적용: 상품 페이지 요청 시 proxy_queue에서 사용 가능한 IP를 하나 가져와 requests 모듈의 proxies 파라미터에 설정하여 요청을 보냅니다.
  4. 프록시 IP 재활용: 성공적으로 요청을 처리한 프록시 IP는 다시 proxy_queue의 끝에 넣어 재사용될 수 있도록 합니다. 실패한 IP는 일정 시간 동안 사용하지 않거나, 유효성 검사를 통해 제거하는 로직을 추가할 수 있습니다.
import requests
import queue
import threading
import time

class ProxyManager:
    def __init__(self, proxy_api_url):
        self.proxy_api_url = proxy_api_url # 프록시 IP를 가져올 API/URL
        self.proxy_queue = queue.Queue()
        self.stop_event = threading.Event()
        self.fetch_thread = threading.Thread(target=self._fetch_proxies_loop, daemon=True)
        self.fetch_thread.start()

    def _fetch_proxies_loop(self):
        while not self.stop_event.is_set():
            if self.proxy_queue.qsize() < 5: # 큐에 프록시가 너무 적으면 새로 가져옴
                print("Fetching new proxies...")
                try:
                    # 실제 프록시 제공 API 호출 (예시)
                    response = requests.get(self.proxy_api_url, timeout=5)
                    response.raise_for_status()
                    # 응답에서 프록시 IP 목록 파싱 (예: ["http://ip:port", "http://ip2:port2"])
                    # 실제 API 응답 형식에 맞게 파싱 로직을 수정해야 합니다.
                    new_proxies_raw = response.json().get("proxies", []) 
                    for proxy_str in new_proxies_raw:
                        self.proxy_queue.put({"http": proxy_str, "https": proxy_str})
                    print(f"Fetched {len(new_proxies_raw)} new proxies. Total: {self.proxy_queue.qsize()}")
                except Exception as e:
                    print(f"Error fetching proxies: {e}")
            time.sleep(30) # 30초마다 확인

    def get_proxy(self):
        try:
            proxy = self.proxy_queue.get(timeout=5)
            print(f"Using proxy: {proxy['http']}")
            return proxy
        except queue.Empty:
            print("No proxies available, waiting...")
            return None

    def put_proxy_back(self, proxy):
        if proxy:
            self.proxy_queue.put(proxy) # 사용한 프록시 재활용

    def stop(self):
        self.stop_event.set()
        if self.fetch_thread.is_alive():
            self.fetch_thread.join(timeout=1) # 스레드가 종료될 때까지 최대 1초 대기

# 스크래퍼 통합 예시 (주석 처리됨):
# proxy_manager = ProxyManager("http://your_proxy_provider_api.com/get_proxies")
# try:
#     proxies = proxy_manager.get_proxy()
#     if proxies:
#         response = requests.get("https://www.amazon.com", proxies=proxies, timeout=10)
#         print(f"Response status with proxy: {response.status_code}")
#         proxy_manager.put_proxy_back(proxies) # 성공 시 재활용
#     else:
#         print("Could not get a proxy.")
# except requests.exceptions.RequestException as e:
#     print(f"Request failed with proxy: {e}")
# except Exception as e:
#     print(f"Other error: {e}")
# finally:
#     proxy_manager.stop()

프록시 IP를 사용할 경우, 요청 헤더와 더불어 proxies 설정을 통해 트래픽을 프록시 서버로 우회시킬 수 있습니다. 이를 통해 웹사이트의 안티-크롤링 시스템을 효과적으로 우회하고 안정적인 데이터 수집을 지속할 수 있습니다.

태그: python 웹 크롤링 아마존 다중 스레딩 데이터 수집

10월 2일 15:41에 게시됨