Quay lại blog

Cách Xây Dựng Một Dòng Thời Gian Web Scraping Được Lên Lịch Với Airflow và Scrapeless

Daniel Kim
Daniel Kim

Lead Scraping Automation Engineer

10-Sep-2026

TL;DR:

  • Một quy trình scraping trong Airflow gồm ba tác vụ — fetch, extract, store — được kết nối bằng cách truyền giá trị trả về, và toàn bộ quá trình ở đây hoàn thành trong 25.2 giây với state=success.
  • Airflow 3 đã đổi tên tham số lập lịch. schedule_interval= nâng TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval', vì vậy hầu hết các hướng dẫn viết cho Airflow 2 dừng lại ở định nghĩa DAG.
  • Cài đặt Airflow với tệp giới hạn của nó. Nếu không có, pip install apache-airflow==3.1.2 đã mất hơn 13 phút trong việc giải quyết phụ thuộc ở đây; với tệp giới hạn, nó đã hoàn thành.
  • Bạn không cần Docker để chạy một DAG. airflow dags test thực thi toàn bộ đồ thị trong quá trình chống lại SQLite.
  • Các giá trị được truyền giữa các tác vụ được tuần tự hóa vào cơ sở dữ liệu siêu dữ liệu — HTML của tác vụ fetch đã được lưu lại ở đó với kích thước 53,889 bytes so với 2,154 cho các bản ghi đã phân tích.
  • Chạy cùng một DAG hai lần để lại 20 hàng thay vì 40, vì ghi là một upsert dựa trên tiêu đề sản phẩm.
  • Bắt đầu lập lịch dựa trên các trang đã được render trên kế hoạch miễn phí Scrapeless.

Một scraper chạy khi bạn nhớ chạy nó chỉ là một script. Biến nó thành một thứ gì đó sản xuất bảng có ngày mỗi sáng có nghĩa là quyết định cách một lần chạy thứ hai ảnh hưởng đến các hàng của ngày hôm qua và nơi chứa thông tin xác thực. Airflow xử lý việc lập lịch và kết nối các tác vụ; hai quyết định đó vẫn thuộc về bạn.

Pipeline dưới đây thu thập một danh mục sách công cộng theo lịch trình hàng ngày: một fetch được render thông qua Universal Scraping API, một phân tích thành các bản ghi, và một upsert vào SQLite. Mỗi số đều được lấy từ việc chạy nó trên Airflow 3.1.2.

Pipeline at a Glance

Giai đoạn Tác vụ Đầu vào Đầu ra Kết quả đã xác minh
Fetch fetch URL danh mục HTML đã render 50,403 ký tự
Extract extract Chuỗi HTML danh sách các bản ghi 20 bản ghi
Store store danh sách bản ghi số lượng hàng 20 hàng trong SQLite

Luồng là fetch → extract → store, và mỗi giai đoạn chuyển giá trị trả về của nó cho giai đoạn tiếp theo. Giữ biên giới tại ba điểm đó cho phép bạn chạy lại một giai đoạn mà không cần lặp lại các giai đoạn khác, và giữ logic phân tích có thể kiểm tra mà không cần cuộc gọi mạng.

Install Airflow Without the Long Resolve

Airflow phụ thuộc vào một đồ thị lớn được cố định, và cài đặt nó mà không có tệp giới hạn tương ứng thì pip sẽ phải quay lại qua các kết hợp phiên bản. Một pip install apache-airflow==3.1.2 đơn giản ở đây vẫn đang giải quyết sau 13 phút; việc cài đặt tương tự với tệp giới hạn đã hoàn thành bình thường.

bash Copy
python3 -m venv .venv
./.venv/bin/pip install "apache-airflow==3.1.2" \
  --constraint "https://raw.githubusercontent.com/apache/airflow/constraints-3.1.2/constraints-3.12.txt"
./.venv/bin/pip install "parsel==1.10.0" "requests==2.32.5"

URL giới hạn mã hóa cả phiên bản Airflow và phiên bản Python, và tham chiếu cài đặt của dự án coi nó như một phần của quá trình cài đặt chứ không phải là một tối ưu hóa. Chỉ AIRFLOW_HOME ở đâu đó rõ ràng và tạo cơ sở dữ liệu siêu dữ liệu:

bash Copy
export AIRFLOW_HOME="$PWD/home"
./.venv/bin/airflow db migrate
text Copy
Airflow database tables created
DB: sqlite:////root/verify-airflow/home/airflow.db
Database migrating done!

SQLite là backend mặc định và nó đủ cho một pipeline trên một máy. Bộ lập lịch cần một cơ sở dữ liệu đồng thời để chạy các tác vụ song song, nhưng không có gì dưới đây yêu cầu điều đó.

Stage 1: Fetch a Rendered Page

Giai đoạn đầu tiên là giai đoạn quyết định liệu phần còn lại của pipeline có thấy dữ liệu hay không. Một máy khách HTTP đơn giản trả về bất kỳ điều gì mà máy chủ gửi trước khi JavaScript chạy, vì vậy tác vụ fetch gọi Universal Scraping API, cái trả về tài liệu đã được render.

python Copy
@task
def fetch() -> str:
    payload = json.dumps({
        "actor": "unlocker.webunlocker",
        "input": {"url": CATEGORY_URL, "js_render": True, "headless": False},
    }).encode()
    request = urllib.request.Request(
        "https://api.scrapeless.com/api/v2/unlocker/request",
        data=payload,
        headers={"Content-Type": "application/json",
                 "x-api-token": os.environ["SCRAPELESS_API_KEY"]},
    )
    with urllib.request.urlopen(request, timeout=180) as response:
        body = json.loads(response.read())
    html = body["data"]
    print(f"fetched {len(html)} chars")
    return html
text Copy
fetched 50403 chars

Bao bì phản hồi là {"code": ..., "data": ...}, nơi data mang HTML dưới dạng một chuỗi. Bởi vì API trả về JSON, tài liệu đã đến đã được giải mã sẵn dưới dạng UTF-8 — £ trong mỗi giá không tồn tại mà không cần xử lý mã hóa trong tác vụ.

Mã API đến từ môi trường thay vì tệp DAG. Airflow đọc môi trường quy trình khi thực thi một tác vụ, vì vậy việc xuất nó trong shell chạy pipeline sẽ giữ thông tin xác thực ra khỏi kho lưu trữ. Airflow Variables và Connections là câu trả lời đầy đủ hơn khi có nhiều hơn một DAG cần nó.

Stage 2: Extract Records

Giai đoạn extract lấy một chuỗi và trả về các bản ghi. Nó không giao tiếp với mạng, điều này có nghĩa là nó có thể được thử nghiệm với một fixture đã lưu bất cứ khi nào markup thay đổi.

python Copy
@task
def extract(html: str) -> list[dict]:
    sel = Selector(text=html)
    rows = []
    for card in sel.css("article.product_pod"):
        rows.append({
            "title": card.css("h3 a::attr(title)").get(),
            "price": card.css("p.price_color::text").get(),
            "rating": (card.css("p.star-rating::attr(class)").get() or "").replace("star-rating", "").strip(),
            "in_stock": "In stock" in (card.css("p.instock.availability::text").getall() or [""])[-1],
        })
    print(f"extracted {len(rows)} records")
    return rows
text Copy
extracted 20 records

Chọn từng thẻ trước và sau đó truy vấn bên trong nó giữ các trường từ cùng một sản phẩm cùng nhau. Việc chọn tất cả các tiêu đề và tất cả các giá như hai danh sách phẳng và ghép chúng lại với nhau sản xuất các cặp không khớp một cách im lặng ngay khi một thẻ thiếu một trường.

Stage 3: Store With an Upsert

Một pipeline lên lịch sẽ ghi lại cùng một hàng một lần nữa vào ngày mai, vì vậy việc ghi phải được định nghĩa dựa trên những gì làm cho một bản ghi là duy nhất thay vì ghi thêm một cách mù quáng.

python Copy
@task
def store(rows: list[dict]) -> int:
    conn = sqlite3.connect(DB_PATH)
    conn.execute("""CREATE TABLE IF NOT EXISTS books (
        title TEXT PRIMARY KEY, price TEXT, rating TEXT,
        in_stock INTEGER, scraped_at TEXT)""")
    now = datetime.utcnow().isoformat(timespec="seconds")
    conn.executemany(
        "INSERT OR REPLACE INTO books VALUES (?,?,?,?,?)",
        [(r["title"], r["price"], r["rating"], int(r["in_stock"]), now) for r in rows],
    )
    conn.commit()
    count = conn.execute("SELECT COUNT(*) FROM books").fetchone()[0]
    conn.close()
    print(f"stored; table now holds {count} rows")
    return count
text Copy
stored; table now holds 20 rows

title là khóa chính, và hành vi INSERT OR REPLACE của SQLite xóa dòng xung đột trước khi chèn dòng mới. Bảng đã lưu sau khi chạy:

text Copy
title                                  price     rating stock
A Murder in Time                       £16.64    One    1
A Study in Scarlet (Sherlock Holmes    £16.73    Two    1
A Time of Torment (Charlie Parker #1   £48.35    Five   1
Boar Island (Anna Pigeon #19)          £59.48    Three  1
Delivering the Truth (Quaker Midwife   £20.89    Four   1
Hide Away (Eve Duncan #20)             £11.84    One    1

Kết Nối Các Giai Đoạn Vào Một DAG

Khai báo DAG mang theo lịch trình; dòng cuối khai báo chuỗi phụ thuộc bằng cách truyền giá trị trả về.

python Copy
@dag(
    dag_id="books_pipeline",
    schedule="0 6 * * *",
    start_date=pendulum.datetime(2026, 9, 1, tz="UTC"),
    catchup=False,
    tags=["scraping"],
)
def books_pipeline():
    ...
    store(extract(fetch()))


books_pipeline()

schedule nhận một biểu thức cron chuẩn, vì vậy 0 6 * * * là 06:00 hàng ngày theo múi giờ của DAG. Thiết lập start_date với một vùng rõ ràng quan trọng hơn nó có vẻ: một DAG gắn chặt vào một vùng quan sát giờ tiết kiệm ánh sáng mặt trời sẽ thay đổi thời gian thực thi của nó hai lần một năm, và các khoảng cách đến từ cơ sở dữ liệu múi giờ IANA. Gắn chặt vào UTC giữ cho khoảng thời gian cố định. catchup=False quan trọng đối với các trình thu thập dữ liệu: với nó được đặt thành True, triển khai một DAG mà start_date là một tháng trước sẽ xếp hàng cho mỗi khoảng thời gian bị lỡ, và một trình thu thập dữ liệu không thể phục hồi một trang như nó đã trông ba tuần trước.

Viết store(extract(fetch())) là toàn bộ đồ thị phụ thuộc. Airflow đọc phần gọi lồng ghép và xây dựng các cạnh.

Chạy Nó Mà Không Cần Docker

Hầu hết các hướng dẫn Airflow bắt đầu với một tệp Docker Compose và một container Postgres. Để phát triển một DAG, airflow dags test chạy toàn bộ đồ thị trong quy trình, mà không cần khởi động bộ lập lịch, máy chủ web, hoặc một container.

bash Copy
export AIRFLOW_HOME="$PWD/home"
export SCRAPELESS_API_KEY="your-api-key"
./.venv/bin/airflow dags test books_pipeline
text Copy
fetched 50403 chars
extracted 20 records
stored; table now holds 20 rows
DagRun Finished: dag_id=books_pipeline, run_duration=25.211182, state=success

Chạy đó là vòng lặp phản hồi nhanh nhất có sẵn trong khi logic phân tích vẫn đang hoạt động. Bộ lập lịch là những gì bạn khởi động khi DAG đang thực hiện điều bạn muốn.

Lập lịch đối với các trang hiển thị ở phía client? Kế hoạch miễn phí Scrapeless bao gồm đủ yêu cầu để chạy một DAG hàng ngày từ đầu đến cuối.

Chạy Hai Lần Không Nên Nhân Đôi Các Dòng Của Bạn

Một đường ống đã được lập lịch chạy không giám sát, vì vậy bài kiểm tra thú vị là lần chạy thứ hai chứ không phải lần đầu tiên. Thực thi cùng một DAG một lần nữa:

text Copy
stored; table now holds 20 rows

Hai mươi, không phải bốn mươi. Upsert có khóa title đã thay thế mỗi hàng và làm mới scraped_at của nó, điều này khiến cho DAG an toàn để kích hoạt thủ công trong khi gỡ lỗi. Một INSERT bình thường sẽ làm cho bảng bị nhân đôi và không để lại cách nào để biết bản sao nào là hiện tại.

Nếu lịch sử quan trọng — theo dõi cách giá cả di chuyển — khóa trở thành cặp tiêu đề và ngày chạy chứ không chỉ là tiêu đề, và hàng của ngày hôm qua vẫn ở lại.

Ranh Giới XCom

Giá trị trả về di chuyển giữa các tác vụ thông qua XCom, và các giá trị XCom được tuần tự hóa vào cơ sở dữ liệu siêu dữ liệu. Ba giai đoạn ở đây đã lưu:

text Copy
task_id      serialized bytes
fetch        53889
extract      2154
store        2

HTML lớn gấp 25 lần kích thước của các bản ghi được phân tích từ nó, và nó nằm trong cơ sở dữ liệu Airflow thay vì trong bộ nhớ. Điều đó là ổn ở kích thước này và sẽ không ổn khi các trang lớn hơn hoặc các lần chạy thường xuyên hơn.

Hình thức có thể mở rộng là phân tích sớm và truyền qua nhỏ: hãy để tác vụ lấy tài liệu thô vào lưu trữ đối tượng và trả về một khóa, sau đó cho phép trích xuất đọc lại. Các ranh giới giai đoạn giữ nguyên và cơ sở dữ liệu siêu dữ liệu vẫn giữ các giá trị có kích thước bản ghi thay vì tài liệu.

Mã Airflow 2 Sẽ Không Chạy Ở Đây Nếu Không Thay Đổi

Hai thay đổi bắt gặp bất kỳ ai theo dõi một hướng dẫn cũ hơn, và chỉ một trong số đó tự thông báo.

python Copy
from airflow.decorators import dag, task   # deprecated
from airflow.sdk import dag, task          # Airflow 3

Đường dẫn nhập cũ vẫn giải quyết và phát ra DeprecatedImportWarning: The airflow.decorators.dag attribute is deprecated. Please use 'airflow.sdk.dag'. Tham số lập lịch là một lỗi nghiêm trọng:

text Copy
TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'

schedule_interval= trở thành schedule=. Vì hầu hết mọi hướng dẫn thu thập dữ liệu Airflow đã công bố đều có trước phiên bản 3, nên TypeError là điều đầu tiên nhiều độc giả gặp phải — và nó xuất hiện vào thời điểm phân tích DAG, trước khi bất kỳ mã thu thập dữ liệu nào chạy.

Kết Luận

Đường ống bao gồm ba hàm và một bộ trang trí: lấy trả về HTML, trích xuất trả về bản ghi, lưu trữ trả về số lượng, và store(extract(fetch())) là đồ thị. Điều khiến nó trở thành một đường ống thay vì một tập lệnh là upsert tồn tại sau lần chạy thứ hai, lịch trình được biểu diễn dưới dạng cron, và ranh giới giai đoạn cho phép bạn phân tích lại mà không cần lấy lại.
Bắt đầu với airflow dags test và SQLite, giữ tài liệu ra khỏi XCom khi chúng phát triển, và gán cài đặt với tệp ràng buộc để điều đầu tiên bạn chiến đấu là markup chứ không phải trình giải quyết phụ thuộc. Đối với phần lưu trữ của cùng một vấn đề, pipeline DuckDB và Parquet của chúng tôi bao gồm đầu ra theo cột, và giá cả liệt kê những gì mà một lịch trình hàng ngày tốn kém trong yêu cầu.

Sẵn sàng để đặt một scraper vào lịch trình? Bắt đầu với kế hoạch miễn phí Scrapeless và chỉ định giai đoạn fetch vào mục tiêu của bạn.

Câu hỏi thường gặp

H: Tôi có cần Docker để chạy Airflow cho việc thu thập dữ liệu không?

Không. airflow dags test <dag_id> chạy toàn bộ đồ thị trong quá trình trái ngược với cơ sở dữ liệu SQLite mặc định, đó là cách mà lần chạy 25.2 giây ở trên đã thực hiện. Docker và Postgres trở nên hữu ích khi bạn muốn trình lập lịch chạy không giám sát với các tác vụ thực thi song song, vì SQLite không hỗ trợ ghi đồng thời mà điều đó yêu cầu.

H: API key nên được đặt ở đâu trong một Airflow DAG?

Không nên ở trong tệp DAG. Đọc nó từ môi trường quá trình giữ cho nó ra khỏi kiểm soát phiên bản, và Airflow Variables hoặc Connections là nhà tốt hơn khi vài DAG cần cùng một thông tin xác thực — cả hai đều được lưu trữ trong cơ sở dữ liệu metadata và được tham chiếu bằng tên thay vì dán vào mã.

H: Làm thế nào để tôi ngăn một scraper theo lịch tạo ra các hàng trùng lặp?

Xác định điều gì làm cho một bản ghi là duy nhất và viết dựa trên khóa đó. Ở đây title là khóa chính và việc ghi là INSERT OR REPLACE, vì vậy lần chạy thứ hai để lại 20 hàng thay vì 40. Nếu bạn muốn lịch sử thay vì trạng thái hiện tại, mở rộng khóa để bao gồm ngày chạy để mỗi bức ảnh của ngày là hàng riêng của nó.

H: Có an toàn để truyền HTML đã thu thập giữa các tác vụ Airflow không?

Ở kích thước nhỏ, có — tài liệu 50,403 ký tự ở trên đã được tuần tự hóa thành 53,889 byte trong cơ sở dữ liệu metadata và chuyển giữa các tác vụ mà không gặp vấn đề gì. Nó không còn là ý tưởng tốt khi tài liệu hoặc tần suất chạy tăng lên, vì mỗi giá trị sống trong cơ sở dữ liệu đó. Ghi tài liệu vào bộ lưu trữ đối tượng và truyền một khóa giữ nguyên ranh giới tác vụ với các giá trị XCom có kích thước bản ghi.

H: Tại sao mã hướng dẫn Airflow của tôi lại gây ra TypeError trên schedule_interval?

Bởi vì nó được viết cho Airflow 2. Phiên bản 3 đã đổi tên tham số thành schedule, và việc truyền tên cũ sẽ gây ra TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' trong khi DAG đang được phân tích. Thay đổi kèm theo là việc nhập: airflow.decorators vẫn hoạt động nhưng cảnh báo, và airflow.sdk là đường dẫn hiện tại.

H: Có nên bật chế độ catchup cho một DAG thu thập dữ liệu không?

Thường thì không. Với catchup=True, một DAG mà start_date của nó ở trong quá khứ sẽ xếp hàng một lần chạy cho mỗi khoảng thời gian bị bỏ lỡ ngay khi nó được triển khai. Hành vi đó phù hợp với việc xử lý lại một tập dữ liệu đã cũ, nhưng một scraper đọc bất cứ điều gì trang hiện giờ hiển thị, vì vậy các lần chạy đó sẽ thu thập dữ liệu của hôm nay và gán nhãn nó bằng ngày lịch sử.

Tại Scrapless, chúng tôi chỉ truy cập dữ liệu có sẵn công khai trong khi tuân thủ nghiêm ngặt các luật, quy định và chính sách bảo mật trang web hiện hành. Nội dung trong blog này chỉ nhằm mục đích trình diễn và không liên quan đến bất kỳ hoạt động bất hợp pháp hoặc vi phạm nào. Chúng tôi không đảm bảo và từ chối mọi trách nhiệm đối với việc sử dụng thông tin từ blog này hoặc các liên kết của bên thứ ba. Trước khi tham gia vào bất kỳ hoạt động cạo nào, hãy tham khảo ý kiến ​​cố vấn pháp lý của bạn và xem xét các điều khoản dịch vụ của trang web mục tiêu hoặc có được các quyền cần thiết.

Bài viết phổ biến nhất

Danh mục