कैसे एयरफ्लो और स्क्रैपलेस के साथ एक निर्धारित वेब स्क्रैपिंग पाइपलाइन बनाएं
Lead Scraping Automation Engineer
TL;DR:
- Airflow में एक स्क्रैपिंग पाइपलाइन तीन कार्यों में होती है — फ़ेच, एक्सट्रैक्ट, स्टोर — जो रिटर्न मूल्यों को पास करके जुड़े होते हैं, और यहां पूर्ण रन 25.2 सेकंड में समाप्त हुआ
state=successके साथ। - Airflow 3 ने शेड्यूलिंग पैरामीटर का नाम बदल दिया।
schedule_interval=TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'को बढ़ाता है, इसलिए Airflow 2 के लिए लिखे गए अधिकांश ट्यूटोरियल DAG परिभाषा पर रुक जाते हैं। - Airflow को इसके प्रतिबंध फ़ाइल के साथ स्थापित करें। इसके बिना,
pip install apache-airflow==3.1.2ने यहां निर्भरता समाधान में 13 मिनट से अधिक समय बिताया; प्रतिबंध फ़ाइल के साथ, यह पूरा हो गया। - DAG चलाने के लिए आपको डॉकर की आवश्यकता नहीं है।
airflow dags testसमग्र ग्राफ को SQLite के खिलाफ प्रक्रिया में निष्पादित करता है। - कार्यों के बीच पास किए गए मानों को मेटाडेटा डेटाबेस में सीरियलाइज़ किया गया है — फ़ेच कार्य का HTML वहां 53,889 बाइट्स में लैंड हुआ है जबकि पार्स किए गए रिकॉर्ड के लिए 2,154 बाइट्स।
- वही DAG दो बार चलाने से 40 की बजाय 20 पंक्तियाँ रह गईं, क्योंकि लेखन उत्पाद शीर्षक पर कुंजीबद्ध एक अपसर्ट है।
- Scrapeless मुफ़्त योजना पर रेंडर की गई पृष्ठों के खिलाफ शेड्यूलिंग शुरू करें।
एक स्क्रैपर जो तब चलता है जब आप इसे चलाने के लिए याद करते हैं, वह एक स्क्रिप्ट है। इसे हर सुबह एक दिनांकित तालिका उत्पन्न करने के लिए एक ऐसी चीज़ में बदलना मतलब यह तय करना है कि एक दूसरी रन कल के पंक्तियों पर क्या करती है और क्रेडेंशियल कहाँ रहते हैं। Airflow शेड्यूलिंग और कार्य वायरिंग को संभालता है; ये दो निर्णय आपके हैं।
नीचे दी गई पाइपलाइन दैनिक रूप से एक सार्वजनिक पुस्तक कैटलॉग एकत्र करती है: यूनिवर्सल स्क्रैपिंग एपीआई के माध्यम से एक रेंडर्ड फ़ेच, रिकॉर्ड्स में पार्स और SQLite में एक अपसर्ट। हर संख्या Airflow 3.1.2 पर इसे चलाने से आती है।
Pipeline at a Glance
| Stage | Task | Input | Output | Verified result |
|---|---|---|---|---|
| Fetch | fetch |
श्रेणी URL | रेंडर्ड HTML | 50,403 अक्षर |
| Extract | extract |
HTML स्ट्रिंग | रिकॉर्ड्स की सूची | 20 रिकॉर्ड |
| Store | store |
रिकॉर्ड सूची | पंक्ति संख्या | SQLite में 20 पंक्तियाँ |
फ्लो है fetch → extract → store, और प्रत्येक चरण अपना रिटर्न मूल्य अगले को सौंपता है। इन तीन बिंदुओं पर सीमाएँ बनाए रखना यही है जो आपको एक चरण को फिर से चलाने की अनुमति देता है बिना अन्य को दोहराए, और यह पार्स लॉजिक को बिना नेटवर्क कॉल के परीक्षण योग्य रखता है।
Install Airflow Without the Long Resolve
Airflow एक बड़े पिन किए गए ग्राफ पर निर्भर करता है, और इसे मिलान करने वाली प्रतिबंध फ़ाइल के बिना स्थापित करने पर pip संस्करण संयोजनों के माध्यम से वापस ट्रैक करता है। यहां एक साधारण pip install apache-airflow==3.1.2 अभी भी 13 मिनट बाद समाधान कर रहा था; प्रतिबंध फ़ाइल के साथ वही इंस्टॉलेशन सामान्य रूप से समाप्त हो गया।
bash
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 दोनों Airflow संस्करण और Python संस्करण को एन्कोड करता है, और परियोजना की स्थापना संदर्भ इसे स्थापना का हिस्सा मानता है न कि एक ऑप्टिमाइजेशन। बिंदु AIRFLOW_HOME को कहीं स्पष्ट रूप से रखें और मेटाडेटा डेटाबेस बनाएं:
bash
export AIRFLOW_HOME="$PWD/home"
./.venv/bin/airflow db migrate
text
Airflow database tables created
DB: sqlite:////root/verify-airflow/home/airflow.db
Database migrating done!
SQLite डिफ़ॉल्ट बैकएंड है और यह एकल-मशीन पाइपलाइन के लिए पर्याप्त है। शेड्यूलर को कार्यों को समानांतर में चलाने के लिए एक समवर्ती डेटाबेस की आवश्यकता होती है, लेकिन यहां नीचे कुछ भी इसकी आवश्यकता नहीं है।
Stage 1: Fetch a Rendered Page
पहला चरण वह है जो यह तय करता है कि पाइपलाइन का शेष हिस्सा डेटा देखता है या नहीं। एक साधारण HTTP क्लाइंट वह सब लौटाता है जो सर्वर भेजता है जब तक जावास्क्रिप्ट नहीं चलती, इसलिए फ़ेच कार्य यूनिवर्सल स्क्रैपिंग एपीआई को कॉल करता है, जो रेंडर्ड दस्तावेज़ लौटाता है।
python
@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
fetched 50403 chars
प्रतिक्रिया लिफाफा {"code": ..., "data": ...} है, जहां data HTML को एक स्ट्रिंग के रूप में ले जाती है। चूंकि एपीआई JSON लौटाता है, दस्तावेज़ पहले से ही UTF-8 के रूप में डिकोडेड आया है — हर कीमत में £ बिना किसी एन्कोडिंग हैंडलिंग के कार्य में जीवित रहता है।
एपीआई टोकन पर्यावरण से आता है न कि DAG फ़ाइल से। Airflow जब कार्य निष्पादित करता है तो प्रक्रिया के पर्यावरण को पढ़ता है, इसलिए पाइपलाइन चलाने वाली शेल में इसे निर्यात करना क्रेडेंशियल को रिपॉजिटरी से बाहर रखता है। Airflow Variables और Connections तब पूरी तरह से उत्तर होती हैं जब एक से अधिक DAG की आवश्यकता होती है।
Stage 2: Extract Records
एक्सट्रैक्ट चरण एक स्ट्रिंग लेता है और रिकॉर्ड लौटाता है। यह कोई नेटवर्क टच नहीं करता, जिसका मतलब है कि इसे जब कभी भी मार्कअप बदलता है तो एक सहेजे गए फिक्स्चर के खिलाफ परखा जा सकता है।
python
@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
extracted 20 records
हर कार्ड को पहले चुनना और फिर इसके अंदर क्वेरी करना एक जैसे उत्पाद के क्षेत्रों को एक साथ रखता है। सभी शीर्षकों और सभी कीमतों को दो सपाट सूचियों के रूप में चुनना और उन्हें zip करना उस क्षण पर चुपचाप मेल खाते जो एक कार्ड में एक क्षेत्र गायब हो।
Stage 3: Store With an Upsert
एक शेड्यूल की गई पाइपलाइन कल वही पंक्तियाँ फिर से लिखती है, इसलिए लेखन को यह परिभाषित करना होता है कि एक रिकॉर्ड को अनूठा बनाने के लिए क्या आवश्यक है न कि अंधाधुंध जोड़ना।
python
@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
stored; table now holds 20 rows
title प्राथमिक कुंजी है, और SQLite का INSERT OR REPLACE व्यवहार टकराने वाली पंक्ति को हटाता है इससे पहले कि नई पंक्ति डाली जाए। रन के बाद संग्रहीत तालिका:
text
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
चरणों को DAG में तार करना
DAG घोषणा कार्यक्रम को ले जाती है; अंतिम पंक्ति लौटाने वाले मानों को पारित करके निर्भरता श्रृंखला की घोषणा करती है।
python
@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 एक मानक क्रोन अभिव्यक्ति लेता है, इसलिए 0 6 * * * DAG के समय क्षेत्र में प्रतिदिन 06:00 है। start_date को स्पष्ट क्षेत्र के साथ सेट करना उतना ही महत्वपूर्ण है जितना यह दिखता है: एक DAG जो एक क्षेत्र पर स्थिर है जो दिन के उजाले की बचत का पालन करता है, साल में दो बार इसके वास्तविक निष्पादन समय को बदलता है, और ऑफसेट IANA समय क्षेत्र डेटाबेस से आते हैं। UTC पर स्थिर होना अंतराल को निश्चित रखता है। catchup=False विशेष रूप से स्क्रैपर के लिए महत्वपूर्ण है: इसे True पर सेट करने के साथ, ऐसा DAG तैनात करना जिसका start_date एक महीने पीछे है, हर चूके हुए अंतराल के लिए एक रन को कतार में लगाता है, और एक स्क्रैपर किसी पृष्ठ को पुनः प्राप्त नहीं कर सकता जैसा कि यह तीन सप्ताह पहले दिखाई दिया।
store(extract(fetch())) लिखना पूरे निर्भरता ग्राफ का है। एयरफ़्लो कॉल नेस्टिंग को पढ़ता है और किनारों का निर्माण करता है।
इसे बिना डॉकर के चलाएँ
अधिकांश एयरफ़्लो वॉटकु शुरू होता है एक डॉकर कंपोज़ फ़ाइल और एक पोस्टग्रेएस कंटेनर के साथ। एक DAG विकसित करने के लिए, airflow dags test पूरे ग्राफ़ को प्रक्रिया में चलाता है, शेड्युलर, वेब सर्वर, या एक कंटेनर शुरू किए बिना।
bash
export AIRFLOW_HOME="$PWD/home"
export SCRAPELESS_API_KEY="your-api-key"
./.venv/bin/airflow dags test books_pipeline
text
fetched 50403 chars
extracted 20 records
stored; table now holds 20 rows
DagRun Finished: dag_id=books_pipeline, run_duration=25.211182, state=success
वह रन सबसे तेज़ फीडबैक लूप है जो उपलब्ध है जबकि पार्स लॉजिक अभी भी चल रहा है। शेड्यूलर वह है जिसे आप तब शुरू करते हैं जब DAG वह कर रहा हो जो आप चाहते हैं।
क्या क्लाइंट-साइड पर रेंडर होने वाले पृष्ठों के खिलाफ शेड्यूल किया गया है? Scrapeless मुफ्त योजना इतनी पर्याप्त अनुरोधों को कवर करता है कि एक दैनिक DAG शुरू से अंत तक चल सके।
दो बार चलाना आपकी पंक्तियों को दोगुना नहीं करना चाहिए
एक शेड्यूल की गई पाइपलाइन बिना देखभाल के चलती है, इसलिए दिलचस्प परीक्षण दूसरा रन है न कि पहला। उसी DAG को फिर से निष्पादित करना:
text
stored; table now holds 20 rows
बीस, चालीस नहीं। title पर की गई अपसर्ट ने प्रत्येक पंक्ति को प्रतिस्थापित किया और इसके scraped_at को पुनः प्रस्तुत किया, जो मैन्युअल रूप से डिबग करते समय DAG को ट्रिगर करने के लिए सुरक्षित बनाता है। एक साधारण INSERT तालिका को दोगुना कर देता और यह बताने का कोई तरीका नहीं छोड़ता कि कौन सी प्रति वर्तमान थी।
यदि इतिहास मायने रखता है — यह ट्रैक करना कि एक मूल्य कैसे गतिशील होता है — कुंजी शीर्षक और रन तिथि का जोड़ा बन जाता है न कि केवल शीर्षक, और कल की पंक्ति बनी रहती है।
XCom सीमा
रिटर्न मान कार्यों के माध्यम से XCom के माध्यम से चलते हैं, और XCom मानों को मेटाडाटा डेटाबेस में श्रृंखलाबद्ध किया जाता है। यहां तीन चरणों को संग्रहीत किया गया:
text
task_id serialized bytes
fetch 53889
extract 2154
store 2
एचटीएमएल उसके द्वारा पार्स किए गए रिकॉर्ड के आकार का 25 गुना है, और यह एयरफ्लो डेटाबेस में sits है न कि मेमोरी में। यह इस आकार पर ठीक है और पृष्ठों के बड़े होने या रन के अधिक बार होने पर ठीक रहना बंद कर देता है।
जो आकार स्केल करता है वह जल्दी पार्स करना और छोटा पास करना है: फ़ेच कार्य को कच्चे दस्तावेज़ को वस्तु संग्रहण में लिखने और एक कुंजी लौटाने दें, फिर निकालें इसे पढ़ने दें। चरण सीमाएँ समान रहती हैं और मेटाडेटा डेटाबेस रिकॉर्ड आकार के मानों को दस्तावेज़ों के स्थान पर बनाए रखता है।
एयरफ्लो 2 कोड यहाँ बिना बदले नहीं चलेगा
दो बदलाव किसी को भी पुराने ट्यूटोरियल का पालन करते समय पकड़ते हैं, और उनमें से केवल एक ही खुद को घोषणा करता है।
python
from airflow.decorators import dag, task # deprecated
from airflow.sdk import dag, task # Airflow 3
पुराना आयात पथ अभी भी हल करता है और DeprecatedImportWarning: The airflow.decorators.dag attribute is deprecated. Please use 'airflow.sdk.dag' निकालता है। शेड्यूलिंग पैरामीटर एक कठोर विफलता है:
text
TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval'
schedule_interval= बन गया schedule=। चूंकि लगभग हर प्रकाशित एयरफ्लो स्क्रैपिंग ट्यूटोरियल संस्करण 3 से पहले है, वह TypeError वह पहली चीज है जो कई पाठकों पर पड़ती है — और यह DAG-पार्स समय पर सामने आती है, इससे पहले कि कोई स्क्रैपिंग कोड चले।
निष्कर्ष
पाइपलाइन तीन कार्य और एक डेकोरेटर है: फ़ेच HTML लौटाता है, निकालना रिकॉर्ड लौटाता है, स्टोर एक गिनती लौटाता है, और store(extract(fetch())) ग्राफ है। इसे एक पाइपलाइन बनाता है न कि एक स्क्रिप्ट की अपसर्ट जो एक दूसरे रन को सहन करती है, शेड्यूल जो क्रोन के रूप में व्यक्त किया गया है, और चरण की सीमा जो आपको पुनः-पार्स करने की अनुमति देती है बिना पुनः-फेचिंग के।
airflow dags test और SQLite के साथ शुरू करें, दस्तावेज़ों को XCom से बाहर रखें जब वे बड़े हो जाएं, और इंस्टॉल को प्रतिबंध फ़ाइल के साथ पिन करें ताकि पहला काम जो आप करें वह मार्कअप हो न कि निर्भरता समाधानकर्ता। उसी समस्या के संग्रहण अंत के लिए, हमारी DuckDB और Parquet पाइपलाइन स्तंभित आउटपुट को कवर करती है, और मूल्य निर्धारण यह सूचीबद्ध करता है कि एक दैनिक शेड्यूल में अनुरोधों की लागत क्या है।
क्या आप एक स्क्रैपर को शेड्यूल पर रखने के लिए तैयार हैं? Scrapeless मुफ्त योजना से शुरू करें और फ़ेच चरण को अपने लक्षित स्थान पर इंगित करें।
FAQ
Q: क्या मुझे स्क्रैपिंग के लिए Airflow चलाने के लिए Docker की जरूरत है?
नहीं। airflow dags test <dag_id> डिफ़ॉल्ट SQLite डेटाबेस के खिलाफ इन-प्रोसेस में पूरा ग्राफ़ चलाता है, जो ऊपर 25.2-सेकंड रन कैसे किया गया। जब आप शेड्यूलर को बिना देखरेख के चलाना चाहते हैं और कार्यों को समानांतर में निष्पादित करते हैं, तो Docker और Postgres उपयोगी हो जाते हैं, क्योंकि SQLite सहसंबंधित लेखन का समर्थन नहीं करता है।
Q: API कुंजी को Airflow DAG में कहाँ होना चाहिए?
DAG फ़ाइल में नहीं। इसे प्रक्रिया वातावरण से पढ़ना संस्करण नियंत्रण से बाहर रखता है, और एयरफ्लो वेरिएबल या कनेक्शंस तब बेहतर होते हैं जब कई DAGs को एक ही क्रेडेंशियल की आवश्यकता होती है — दोनों मेटाडेटा डेटाबेस में संग्रहीत होते हैं और नाम द्वारा संदर्भित किए जाते हैं न कि कोड में चिपकाए जाते हैं।
Q: मैं एक निर्धारित स्क्रैपर को डुप्लिकेट पंक्तियाँ बनाने से कैसे रोक सकता हूँ?
यह परिभाषित करें कि रिकॉर्ड को अद्वितीय बनाने के लिए क्या आवश्यक है और उस कुंजी के खिलाफ लिखें। यहाँ title प्राथमिक कुंजी है और लिखना INSERT OR REPLACE है, इसलिए दूसरा रन 40 की बजाय 20 पंक्तियाँ छोड़ गया। यदि आप वर्तमान स्थिति के बजाय इतिहास चाहते हैं, तो कुंजी को रन तिथि में शामिल करने के लिए चौड़ा करें ताकि हर दिन का स्नैपशॉट अपनी खुद की पंक्ति हो।
Q: क्या यह सुरक्षित है कि स्क्रैप किए गए HTML को Airflow कार्यों के बीच पास करें?
छोटी आकारों पर, हाँ — ऊपर 50,403-चरित्र दस्तावेज़ ने मेटाडेटा डेटाबेस में 53,889 बाइट्स में सीरियल किया और कार्यों के बीच बिना परेशानी के स्थानांतरित हो गया। जैसे-जैसे दस्तावेज़ या रन की आवृत्ति बढ़ती है, यह एक अच्छा विचार बनना बंद कर देता है, क्योंकि प्रत्येक मान उसी डेटाबेस में रहता है। दस्तावेज़ को वस्तु संग्रहण में लिखना और एक कुंजी को पास करना समान कार्य सीमाओं को बनाए रखता है जिसमें रिकॉर्ड-आकार के XCom मान होते हैं।
Q: क्यों मेरे Airflow ट्यूटोरियल कोड पर schedule_interval पर TypeError उठता है?
क्योंकि इसे Airflow 2 के लिए लिखा गया था। संस्करण 3 ने schedule नामक पैरामीटर का नाम बदल दिया, और पुराने नाम को पास करना TypeError: DAG.__init__() got an unexpected keyword argument 'schedule_interval' को उठाता है जबकि DAG का विश्लेषण किया जा रहा है। सहायक परिवर्तन आयात है: airflow.decorators अभी भी काम करता है लेकिन चेतावनी देता है, और airflow.sdk वर्तमान पथ है।
Q: क्या scraping DAG के लिए catchup सक्षम होना चाहिए?
आमतौर पर नहीं। catchup=True के साथ, एक DAG जिसकी start_date अतीत में है, हर मिस किए गए अंतराल के लिए एक रन को तब तक कतारबद्ध करती है जब तक कि इसे लागू नहीं किया जाता। यह व्यवहार एक डेटेड डेटा सेट को पुनः-प्रसंस्करण के लिए उपयुक्त है, लेकिन एक स्क्रैपर जो भी पृष्ठ वर्तमान में दिखाता है उसे पढ़ता है, इसलिए वे रन आज का डेटा एकत्र करते हैं और इसे ऐतिहासिक तिथियों के साथ लेबल करते हैं।
स्क्रैपलेस में, हम केवल सार्वजनिक रूप से उपलब्ध डेटा का उपयोग करते हैं, जबकि लागू कानूनों, विनियमों और वेबसाइट गोपनीयता नीतियों का सख्ती से अनुपालन करते हैं। इस ब्लॉग में सामग्री केवल प्रदर्शन उद्देश्यों के लिए है और इसमें कोई अवैध या उल्लंघन करने वाली गतिविधियों को शामिल नहीं किया गया है। हम इस ब्लॉग या तृतीय-पक्ष लिंक से जानकारी के उपयोग के लिए सभी देयता को कोई गारंटी नहीं देते हैं और सभी देयता का खुलासा करते हैं। किसी भी स्क्रैपिंग गतिविधियों में संलग्न होने से पहले, अपने कानूनी सलाहकार से परामर्श करें और लक्ष्य वेबसाइट की सेवा की शर्तों की समीक्षा करें या आवश्यक अनुमतियाँ प्राप्त करें।



