רכיבי Dataflow

רכיבי Dataflow מאפשרים לשלוח משימות של Apache Beam ל-Dataflow לביצוע. ב-Dataflow, משאב Job מייצג משימת Dataflow.

‫ Google Cloud SDK כולל את האופרטורים הבאים ליצירת משאבי Job ולמעקב אחר ההפעלה שלהם:

בנוסף, ה-SDK‏ כולל את הרכיב WaitGcpResourcesOp, שבעזרתו אפשר לצמצם את העלויות בזמן הפעלת משימות Dataflow. Google Cloud

DataflowFlexTemplateJobOp

האופרטור DataflowFlexTemplateJobOp מאפשר ליצור רכיב של Gemini Enterprise Agent Platform Pipelines כדי להפעיל Dataflow Flex Template.

ב-Dataflow, משאב LaunchFlexTemplateParameter מייצג תבנית Flex להפעלה. הרכיב הזה יוצר משאב LaunchFlexTemplateParameter ואז שולח בקשה ל-Dataflow ליצור משימה על ידי הפעלת התבנית. אם התבנית מופעלת בהצלחה, Dataflow מחזיר משאב Job.

רכיב Dataflow Flex Template מסתיים כשמתקבל משאב Job מ-Dataflow. הרכיב מוציא job_id כ-serialized gcp_resources proto. אפשר להעביר את הפרמטר הזה לרכיב WaitGcpResourcesOp כדי להמתין לסיום המשימה ב-Dataflow.

DataflowPythonJobOp

האופרטור DataflowPythonJobOp מאפשר ליצור רכיב של צינורות עיבוד נתונים בפלטפורמת הסוכנים של Gemini Enterprise שמכין נתונים על ידי שליחת משימת Apache Beam מבוססת-Python ל-Dataflow לביצוע.

קוד ה-Python של משימת Apache Beam פועל באמצעות Dataflow Runner. כשמריצים את צינור הנתונים באמצעות שירות Dataflow, ה-Runner מעלה את הקוד הניתן להפעלה למיקום שצוין בפרמטר python_module_path ואת התלות לקטגוריה של Cloud Storage (שצוינה בפרמטר temp_location), ואז יוצר משימת Dataflow שמבצעת את צינור הנתונים של Apache Beam במשאבים מנוהלים ב- Google Cloud.

מידע נוסף על Dataflow Runner זמין במאמר בנושא שימוש ב-Dataflow Runner.

רכיב Python של Dataflow מקבל רשימה של ארגומנטים שמועברים באמצעות Beam Runner לקוד Apache Beam. הארגומנטים האלה מוגדרים על ידי args. לדוגמה, אפשר להשתמש בארגומנטים האלה כדי להגדיר את apache_beam.options.pipeline_options ולציין רשת, תת-רשת, מפתח הצפנה בניהול הלקוח (CMEK) ואפשרויות אחרות כשמריצים משימות Dataflow.

WaitGcpResourcesOp

לעתים קרובות, לוקח הרבה זמן להשלים משימות Dataflow. העלויות של מאגר busy-wait (המאגר שמפעיל את משימת Dataflow וממתין לתוצאה) יכולות להיות גבוהות.

אחרי ששולחים את עבודת Dataflow באמצעות Beam runner, הרכיב DataflowPythonJobOp מסתיים מיד ומחזיר פרמטר פלט job_id בתור proto gcp_resources מסוג serialized. אפשר להעביר את הפרמטר הזה לרכיב WaitGcpResourcesOp כדי להמתין לסיום של משימת Dataflow.

    dataflow_python_op = DataflowPythonJobOp(
        project=project_id,
        location=location,
        python_module_path=python_file_path,
        temp_location = staging_dir,
        requirements_file_path = requirements_file_path,
        args = ['--output', OUTPUT_FILE],
    )
  
    dataflow_wait_op =  WaitGcpResourcesOp(
        gcp_resources = dataflow_python_op.outputs["gcp_resources"]
    )

‫Gemini Enterprise Agent Platform Pipelines מבצע אופטימיזציה של WaitGcpResourcesOp כדי להריץ אותו באופן חסר שרתים, ללא עלות.

אם הרכיבים DataflowPythonJobOp ו-DataflowFlexTemplateJobOp לא עונים על הדרישות שלכם, אתם יכולים גם ליצור רכיב משלכם שמפיק את הפרמטר gcp_resources ומעביר אותו לרכיב WaitGcpResourcesOp.

מידע נוסף על יצירת פרמטר פלט gcp_resources זמין במאמר כתיבת רכיב להצגת קישור למסוף Google Cloud .

הפניית API

מדריכים

היסטוריית גרסאות ונתוני גרסה

מידע נוסף על היסטוריית הגרסאות והשינויים ב- Google Cloud Pipeline Components SDK זמין בהערות לגבי הגרסה של Pipeline Components SDK.Google Cloud

אנשי קשר לתמיכה טכנית

אם יש לכם שאלות, אפשר לפנות אל kubeflow-pipelines-components@google.com.