स्ट्रीमिंग और GPU इंजनों के साथ काम करना

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

Liam Brannigan

Data Scientist & Polars Contributor

क्वेरी एक्जीक्यूशन इंजन

API से क्वेरी ऑप्टिमाइज़र होकर इन-मेमाेरी इंजन तक जाने वाली क्वेरी का फ्लो डायग्राम.

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

क्वेरी एक्जीक्यूशन इंजन

API से क्वेरी ऑप्टिमाइज़र होकर इन-मेमाेरी इंजन तक जाने वाली क्वेरी का फ्लो डायग्राम.

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

क्वेरी एक्जीक्यूशन इंजन

API से क्वेरी ऑप्टिमाइज़र होकर इन-मेमाेरी और स्ट्रीमिंग इंजनों तक जाने वाली क्वेरी का फ्लो डायग्राम.

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

क्वेरी एक्जीक्यूशन इंजन

API से क्वेरी ऑप्टिमाइज़र होकर इन-मेमाेरी, स्ट्रीमिंग और GPU इंजनों तक जाने वाली क्वेरी का फ्लो डायग्राम.

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

एक लेज़ी रिक्वेस्ट सारांश

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

query = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .group_by("DEPARTMENT")
    .len()
    .sort("len", descending=True)
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

डिफॉल्ट इंजन का उपयोग

query.collect()
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

स्ट्रीमिंग इंजन को टार्गेट करना

query.collect(
    engine="streaming"
)
shape: (5, 2)
| DEPARTMENT        | len     |
| ---               | ---     |
| str               | u32     |
|-------------------|---------|
| 311 City Services | 4859161 |
| Sanitation        | 3406631 |
| Aviation          | 2337842 |
| Transportation    | 1468240 |
| Water             | 416535  |
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

स्ट्रीमिंग इंजन को टार्गेट करना

बड़ी टेबल को स्ट्रीम करते हुए छवि.

  • Unsupported स्ट्रीमिंग स्टेप्स अपने-आप डिफॉल्ट इंजन पर वापस चले जाते हैं
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

स्ट्रीमिंग इंजन को टार्गेट करना

import polars as pl
pl.Config.set_engine_affinity("streaming")
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

GPU इंजन को टार्गेट करना

query.collect(
    engine="gpu"
)

$$

  • NVIDIA GPU हार्डवेयर आवश्यक है
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

इंजन चुनना

$$

  • पहले डिफॉल्ट इंजन; केवल मेमोरी की चिंता हो तो स्विच करें

$$

  • बड़े, out-of-memory क्वेरी के लिए स्ट्रीमिंग

$$

  • स्पीड के लिए GPU, पर तभी जब NVIDIA हार्डवेयर उपलब्ध हो

इंजन चुनना

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

बैच में काम करना

completed_requests = (
    requests
    .filter(pl.col("STATUS") == "Completed")
    .select("SR_NUMBER", "TYPE", "DEPARTMENT", "CREATED_DATE")
)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

प्रत्येक बैच प्रोसेस करना

for batch in completed_requests.collect_batches(

):


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

प्रत्येक बैच प्रोसेस करना

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):


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

प्रत्येक बैच प्रोसेस करना

for batch in completed_requests.collect_batches(
    chunk_size=50_000,
):
    print(batch.shape)
    batch_json = batch.write_json()
(50000, 4)
(50000, 4)
...
(50000, 4)
(12409, 4)
Polars के साथ Data Pipelines का स्केलिंग और ऑप्टिमाइज़ेशन

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

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

Preparing Video For Download...