שימוש ב-Ray ב-Agent Platform עם BigQuery

כשמריצים אפליקציית Ray בפלטפורמת הסוכנים של Gemini Enterprise, צריך להשתמש ב-BigQuery כמסד הנתונים בענן. בחלק הזה מוסבר איך לקרוא ממסד נתונים של BigQuery ולכתוב בו מתוך אשכול Ray בפלטפורמת הסוכנים של Gemini Enterprise. ההנחיות בחלק הזה מניחות שאתם משתמשים ב-SDK של פלטפורמת הסוכנים ל-Python.

כדי לקרוא ממערך נתונים ב-BigQuery, יוצרים מערך נתונים חדש ב-BigQuery או משתמשים במערך נתונים קיים.

ייבוא והפעלה של לקוח Ray ב-Agent Platform

אם אתם מחוברים לאשכול Ray ב-Gemini Enterprise Agent Platform, צריך להפעיל מחדש את ליבת המערכת ולהריץ את הקוד הבא. המשתנה runtime_env נדרש בזמן החיבור כדי להריץ פקודות BigQuery.

import ray
from google.cloud import aiplatform

# The CLUSTER_RESOURCE_NAME is the one returned from vertex_ray.create_ray_cluster.
address = 'vertex_ray://{}'.format(CLUSTER_RESOURCE_NAME)

runtime_env = {
    "pip":
       ["google-cloud-aiplatform[ray]","ray==2.47.1"]
  }

ray.init(address=address, runtime_env=runtime_env)

קריאת נתונים מ-BigQuery

קריאת נתונים ממערך הנתונים שלכם ב-BigQuery. משימת Ray צריכה לבצע את פעולת הקריאה.

aiplatform.init(project=PROJECT_ID, location=LOCATION)

@ray.remote
def run_remotely():
    import vertex_ray
    dataset = DATASET
    parallelism = PARALLELISM
    query = QUERY

    ds = vertex_ray.data.read_bigquery(
        dataset=dataset,
        parallelism=parallelism,
        query=query
    )
    ds.materialize()

כאשר:

  • PROJECT_ID: Google Cloud מזהה הפרויקט. מזהה הפרויקט מופיע בדף welcome במסוף Google Cloud .

  • LOCATION: המיקום שבו מאוחסן Dataset. לדוגמה, us-central1.

  • DATASET: מערך נתונים ב-BigQuery. הפורמט הנדרש הוא dataset.table. מגדירים את הערך None אם מספקים שאילתה.

  • PARALLELISM: מספר שלם שמשפיע על מספר משימות הקריאה שנוצרות במקביל. יכול להיות שייווצרו פחות זרמי קריאה ממה שביקשתם.

  • QUERY: מחרוזת שמכילה שאילתת SQL לקריאה ממסד נתונים של BigQuery. מגדירים את הערך None אם לא נדרשת שאילתה.

טרנספורמציה של נתונים

אפשר לעדכן ולמחוק שורות ועמודות מטבלאות BigQuery באמצעות pyarrow או pandas. אם רוצים להשתמש בהמרות של pandas, צריך להשאיר את סוג הקלט כ-pyarrow ולהמיר ל-pandas בתוך פונקציה בהגדרת משתמש (UDF), כדי שתוכלו לזהות שגיאות בהמרת סוגים של pandas בתוך ה-UDF. משימת Ray צריכה לבצע את השינוי.

@ray.remote
def run_remotely():
    # BigQuery Read first
    import pandas as pd
    import pyarrow as pa

    def filter_batch(table: pa.Table) -> pa.Table:
        df = table.to_pandas(types_mapper={pa.int64(): pd.Int64Dtype()}.get)
        # PANDAS_TRANSFORMATIONS_HERE
        return pa.Table.from_pandas(df)

    ds = ds.map_batches(filter_batch, batch_format="pyarrow").random_shuffle()
    ds.materialize()

    # You can repartition before writing to determine the number of write blocks
    ds = ds.repartition(4)
    ds.materialize()

כתיבת נתונים ל-BigQuery

מכניסים נתונים למערך הנתונים ב-BigQuery. משימת Ray צריכה לבצע את הכתיבה.

@ray.remote
def run_remotely():
    # BigQuery Read and optional data transformation first
    dataset=DATASET
    vertex_ray.data.write_bigquery(
        ds,
        dataset=dataset
    )

כאשר:

  • DATASET: מערך נתונים ב-BigQuery. מערך הנתונים צריך להיות בפורמט dataset.table.

המאמרים הבאים