मल्टिफाइल डेटासेट के साथ काम

Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

Liam Brannigan

Data Scientist & Polars Contributor

रोज़ाना 311 फाइलें

311_requests_daily_csv/
  311_requests_2025-12-31.csv
  311_requests_2026-01-01.csv
  311_requests_2026-01-02.csv
  311_requests_2026-01-03.csv
  311_requests_2026-01-04.csv
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

कई CSV फाइलें स्कैन करना

pl.scan_csv(


)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

कई CSV फाइलें स्कैन करना

pl.scan_csv(
    "311_requests_daily_csv/311_requests_2026*.csv",
    try_parse_dates=True,
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

कई CSV फाइलें स्कैन करना

pl.scan_csv(
    "311_requests_daily_csv/311_requests_2026*.csv",
    try_parse_dates=True,
).select("REPORTED_AT", "STATUS", "WARD").collect()
shape: (4, 3)
| REPORTED_AT         | STATUS      | WARD |
| ---                 | ---         | ---  |
| datetime[us]        | str         | i64  |
|---------------------|-------------|------|
| 2026-01-01 07:10:00 | In Progress | 46   |
| 2026-01-01 08:25:00 | Open        | 46   |
| 2026-01-04 07:45:00 | Open        | 12   |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

हाइव पार्टिशन लिखना

pl.read_csv(
    "311_Service_Requests.csv",
    try_parse_dates=True,
).write_parquet(


)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

हाइव पार्टिशन लिखना

pl.read_csv(
    "311_Service_Requests.csv",
    try_parse_dates=True,
).write_parquet(
    "311_requests_parquet/",

)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

हाइव पार्टिशन लिखना

pl.read_csv(
    "311_Service_Requests.csv",
    try_parse_dates=True,
).write_parquet(
    "311_requests_parquet/",
    partition_by=["CREATED_DATE"],
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन किए गए डेटासेट

311_requests_parquet/
  CREATED_DATE=2025-12-31/
    data.parquet


Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन किए गए डेटासेट

311_requests_parquet/
  CREATED_DATE=2025-12-31/
    data.parquet
  CREATED_DATE=2026-01-01/
    data.parquet
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

हाइव-पार्टिशन्ड डेटासेट स्कैन करना

requests = pl.scan_parquet(
    "311_requests_parquet/",
    hive_partitioning=True,
)

$$

  • Partition values → regular columns
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन-अवेयर क्वेरी

requests.filter(
    pl.col("CREATED_DATE") >= date(2026, 1, 1)
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन-अवेयर क्वेरी

requests.filter(
    pl.col("CREATED_DATE") >= date(2026, 1, 1)
).group_by("DEPARTMENT").len().collect()
shape: (2, 2)
| DEPARTMENT     | len |
| ---            | --- |
| str            | u32 |
|----------------|-----|
| Transportation | 208 |
| Sanitation     | 175 |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन-अनअवेयर क्वेरी!

requests.filter(
    pl.col("YEAR") >= 2026
)

$$

  • Polars सीधे partition कॉलम पर फ़िल्टर करता है 💡
  • Derived कॉलम पर फ़िल्टरिंग यह optimization ट्रिगर नहीं करती 💡
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

स्कीमा ड्रिफ्ट की समस्या

# 2026-01-01 daily file
REPORTED_AT,STATUS,WARD,BLOCK_CODE,REPORT_TEXT
2026-01-01 07:10:00,In Progress,46,0021,...
# 2026-01-02 daily file
REPORTED_AT,STATUS,WARD,REPORT_TEXT
2026-01-02 09:05:00,Open,29,...
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

मिसिंग कॉलम जोड़ना

requests = pl.scan_csv(
    "311_requests_daily_csv/311_requests_2026*.csv",
    try_parse_dates=True,
    missing_columns="insert",
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

कम्बाइंड स्कीमा जाँचना

requests.select(
    "REPORTED_AT", "STATUS", "WARD", "BLOCK_CODE"
).collect()
shape: (8, 4)
| REPORTED_AT         | STATUS      | WARD | BLOCK_CODE |
| ---                 | ---         | ---  | ---        |
| datetime[us]        | str         | i64  | i64        |
|---------------------|-------------|------|------------|
| 2026-01-01 07:10:00 | In Progress | 46   | 21         |
| 2026-01-01 08:25:00 | Open        | 46   | 21         |
| 2026-01-02 09:05:00 | Open        | 29   | null       |
| 2026-01-02 10:42:00 | Completed   | 29   | null       |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

dtype मिसमैच की समस्या

# 2025-12-31 daily file
WARD
35

# 2026-01-03 daily file
WARD
"44"
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

अलग dtypes के साथ दो स्कैन

ward_int = pl.scan_csv(
    "311_requests_daily_csv/311_requests_2025-12-31.csv",
    try_parse_dates=True,
)

ward_str = pl.scan_csv(
    "311_requests_daily_csv/311_requests_2026-01-03.csv",
    try_parse_dates=True,
    schema_overrides={"WARD": pl.String},
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

vertical_relaxed से जोड़ना

combined = pl.concat(
    [ward_int, ward_str],
    how="vertical_relaxed",
)

$$

  • Supertype: ऐसा dtype जो दोनों इनपुट के मान रख सके
  • यहाँ String दोनों Int64 और String के मान रख सकता है
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

relaxed concat जाँचना

combined.collect()
shape: (4, 5)
| REPORTED_AT         | STATUS    | WARD | BLOCK_CODE | REPORT_TEXT                      |
| ---                 | ---       | ---  | ---        | ---                              |
| datetime[us]        | str       | str  | i64        | str                              |
|---------------------|-----------|------|------------|----------------------------------|
| 2025-12-31 16:20:00 | Open      | 35   | 17         | Massive pothole getting worse    |
| 2025-12-31 18:05:00 | Open      | 35   | 17         | Cyclists forced into traffic     |
| 2026-01-03 07:35:00 | Open      | 44   | 12         | How is this still not repaired   |
| 2026-01-03 08:15:00 | Completed | 44   | 12         | Every bus hit shakes the vehicle |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

अभ्यास करते हैं!

Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

Preparing Video For Download...