बड़े क्वेरी आउटपुट को डिस्क पर सिंक करना

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

Liam Brannigan

Data Scientist & Polars Contributor

बड़ा आउटपुट समस्या

एक लेज़ी क्वेरी का डायग्राम, जो बड़ा row-preserving एक्सट्रैक्ट बनाती है जिसे सीधे डिस्क पर लिखना चाहिए.

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

क्लीन किया हुआ रिक्वेस्ट डेटासेट

requests = pl.scan_parquet("311_Service_Requests.parquet")

clean_requests = (
    requests





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

क्लीन किया हुआ रिक्वेस्ट डेटासेट

requests = pl.scan_parquet("311_Service_Requests.parquet")

clean_requests = (
    requests
    .with_columns(
        pl.col("CREATED_DATE").dt.date(),
        pl.col("WARD").cast(pl.Int16),
        (pl.col("STATUS") == "Completed").alias("IS_COMPLETED"),
    )
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पहले कलेक्ट, फिर लिखें

clean_requests.collect(
    engine="streaming"
).write_parquet(
    "311_requests_clean.parquet"
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

Parquet पर सिंक करना

clean_requests.sink_parquet(
    "311_requests_clean.parquet",


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

सिंक को नियंत्रित करना

clean_requests.sink_parquet(
    "311_requests_clean.parquet",
    compression="zstd",
    row_group_size=100_000,
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

अन्य सिंक मेथड्स

df.sink_csv("output.csv")
df.sink_ndjson("output.ndjson")
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन किया हुआ आउटपुट

clean_requests.sink_parquet(
    pl.PartitionBy(



    ),

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

पार्टिशन किया हुआ आउटपुट

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),

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

पार्टिशन किया हुआ आउटपुट

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key="CREATED_DATE",
        max_rows_per_file=1_000_000,
    ),
    mkdir=True,
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

पार्टिशन किया हुआ आउटपुट

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

मल्टी-कॉलम पार्टिशन

clean_requests.sink_parquet(
    pl.PartitionBy(
        "311_requests_clean/",
        key=["STATUS", "CREATED_DATE"],
    ),
    mkdir=True,
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

लेज़ी सिंक बनाना

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)





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

लेज़ी सिंक बनाना

internal_sink = clean_requests.sink_parquet(
    "311_requests_internal.parquet",
    lazy=True,
)

public_sink = clean_requests.drop("REQUEST_ID","REPORTER").sink_parquet(
    "311_requests_public.parquet",
    lazy=True,
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

मल्टीप्लेक्सिंग सिंक

pl.collect_all(
    [internal_sink, public_sink],
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

लिखे गए आउटपुट

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

कस्टम बैच सिंक

def send_batch(batch: pl.DataFrame) -> None:
    print(batch.height)
    batch_json = batch.write_json()
    # send batch_json to an API
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

कस्टम बैच सिंक

clean_requests.sink_batches(
    send_batch,
    chunk_size=50_000,
)
50000
50000
...
29268
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

सिंक वर्कफ़्लो

  • बड़े क्वेरी परिणाम सीधे डिस्क पर लिखने के लिए sink_parquet() का उपयोग करें
  • आउटपुट को पार्टिशन किए गए डेटासेट में बाँटने के लिए PartitionBy() का उपयोग करें
  • कई आउटपुट साथ में प्लान और रन करने हेतु collect_all() के साथ लेज़ी सिंक का उपयोग करें
  • जब हर बैच कस्टम फंक्शन को भेजना हो तो sink_batches() का उपयोग करें

डेटा सॉर्टिंग सिस्टम का आइसोमेट्रिक दृश्य, जो डेटा को व्यवस्थित पार्टिशन दराजों में बाँट रहा है

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

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

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

Preparing Video For Download...