מעקב אחרי שאילתות מתמשכות

אתם יכולים לעקוב אחרי שאילתות רציפות ב-BigQuery באמצעות הכלים הבאים של BigQuery:

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

שימוש בתצוגות INFORMATION_SCHEMA

אתם יכולים להשתמש במספר תצוגות INFORMATION_SCHEMA כדי לעקוב אחרי שאילתות מתמשכות והזמנות של שאילתות מתמשכות.

צפייה בפרטי המשרה

אפשר להשתמש בתצוגה JOBS כדי לקבל מטא-נתונים של משימות של שאילתות מתמשכות.

השאילתה הבאה מחזירה את המטא-נתונים של כל השאילתות המתמשכות הפעילות. המטא-נתונים כוללים את חותמת הזמן של סימן המים של הפלט, שמייצגת את הנקודה שעד אליה השאילתה המתמשכת עיבדה נתונים בהצלחה.

  1. במסוף Cloud de Confiance , עוברים לדף BigQuery.

    כניסה ל-BigQuery

  2. בעורך השאילתות, מריצים את השאילתה הבאה:

    SELECT
      start_time,
      job_id,
      user_email,
      query,
      state,
      reservation_id,
      continuous_query_info.output_watermark
    FROM `PROJECT_ID.region-REGION.INFORMATION_SCHEMA.JOBS`
    WHERE
      creation_time > TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 7 day)
      AND continuous IS TRUE
      AND state = "RUNNING"
    ORDER BY
      start_time DESC

    מחליפים את מה שכתוב בשדות הבאים:

צפייה בפרטי ההקצאה של ההזמנה

אפשר להשתמש בתצוגות ASSIGNMENTS ו-RESERVATIONS כדי לקבל פרטים על הקצאת מקום שמור של שאילתות מתמשכות.

החזרת פרטי הקצאת הזמנות לשאילתות מתמשכות:

  1. במסוף Cloud de Confiance , עוברים לדף BigQuery.

    כניסה ל-BigQuery

  2. בעורך השאילתות, מריצים את השאילתה הבאה:

    SELECT
      reservation.reservation_name,
      reservation.slot_capacity
    FROM
      `ADMIN_PROJECT_ID.region-LOCATION.INFORMATION_SCHEMA.ASSIGNMENTS`
        AS assignment
    INNER JOIN
      `ADMIN_PROJECT_ID.region-LOCATION.INFORMATION_SCHEMA.RESERVATIONS`
        AS reservation
      ON (assignment.reservation_name = reservation.reservation_name)
    WHERE
      assignment.assignee_id = 'PROJECT_ID'
      AND job_type = 'CONTINUOUS';

    מחליפים את מה שכתוב בשדות הבאים:

    • ADMIN_PROJECT_ID: המזהה של פרויקט הניהול שבבעלותו נמצאת ההזמנה.
    • LOCATION: המיקום של ההזמנה.
    • PROJECT_ID: המזהה של הפרויקט שמוקצה להזמנה. מוחזר רק מידע על שאילתות רציפות שפועלות בפרויקט הזה.

צפייה במידע על צריכת משבצות

אפשר להשתמש בתצוגות ASSIGNMENTS,‏ RESERVATIONS ו-JOBS_TIMELINE כדי לקבל מידע על צריכת יחידות קיבולת (Slot) של שאילתות מתמשכות.

מידע על ניצול משבצות להחזרת נתונים בשאילתות מתמשכות:

  1. במסוף Cloud de Confiance , עוברים לדף BigQuery.

    כניסה ל-BigQuery

  2. בעורך השאילתות, מריצים את השאילתה הבאה:

    SELECT
      jobs.period_start,
      reservation.reservation_name,
      reservation.slot_capacity,
      SUM(jobs.period_slot_ms) / 1000 AS consumed_total_slots
    FROM
      `ADMIN_PROJECT_ID.region-LOCATION.INFORMATION_SCHEMA.ASSIGNMENTS`
        AS assignment
    INNER JOIN
      `ADMIN_PROJECT_ID.region-LOCATION.INFORMATION_SCHEMA.RESERVATIONS`
        AS reservation
      ON (assignment.reservation_name = reservation.reservation_name)
    INNER JOIN
      `PROJECT_ID.region-LOCATION.INFORMATION_SCHEMA.JOBS_TIMELINE` AS jobs
      ON (
        UPPER(CONCAT('ADMIN_PROJECT_ID:LOCATION.', assignment.reservation_name))
        = UPPER(jobs.reservation_id))
    WHERE
      assignment.assignee_id = 'PROJECT_ID'
      AND assignment.job_type = 'CONTINUOUS'
      AND jobs.period_start
        BETWEEN TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 1 DAY)
        AND CURRENT_TIMESTAMP()
    GROUP BY 1, 2, 3
    ORDER BY jobs.period_start DESC;

    מחליפים את מה שכתוב בשדות הבאים:

    • ADMIN_PROJECT_ID: המזהה של פרויקט הניהול שבבעלותו נמצאת ההזמנה.
    • LOCATION: המיקום של ההזמנה.
    • PROJECT_ID: המזהה של הפרויקט שמוקצה להזמנה. מוחזר רק מידע על שאילתות רציפות שפועלות בפרויקט הזה.

אפשר גם לעקוב אחרי הזמנות של שאילתות מתמשכות באמצעות כלים אחרים, כמו Metrics Explorer ותרשימים של משאבים אדמיניסטרטיביים. מידע נוסף זמין במאמר בנושא מעקב אחרי הזמנות ב-BigQuery.

שימוש בתרשים של ביצוע השאילתה

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

צפייה בהיסטוריית הפעולות

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

חשוב לדעת שהרשימה ההיסטורית של המשימות ממוינת לפי שעת ההתחלה של המשימה, ולכן יכול להיות ששאילתות רציפות שפועלות כבר זמן מה לא יופיעו קרוב לתחילת הרשימה.

שימוש בכלי לחיפוש משרות

בכלי לבדיקת משימות, מסננים את המשימות כדי להציג שאילתות מתמשכות. לשם כך, מגדירים את המסנן Job category (קטגוריית משימות) לערך Continuous query (שאילתה מתמשכת).

שימוש ב-Cloud Monitoring

אפשר להשתמש ב-Cloud Monitoring כדי להציג מדדים שספציפיים לשאילתות רציפות ב-BigQuery. מידע נוסף זמין במאמרים יצירה של מרכזי בקרה, תרשימים והתראות ומדדים שזמינים להצגה חזותית.

התראה על שאילתות שנכשלו

במקום לבדוק באופן שגרתי אם השאילתות הרציפות נכשלו, כדאי ליצור התראה שתשלח לכם הודעה על כשל. אחת הדרכים לעשות זאת היא ליצור מדד מותאם אישית מבוסס-יומן ב-Cloud Logging עם מסנן למשרות שלכם, ומדיניות התראות ב-Cloud Monitoring שמבוססת על המדד הזה:

  1. כשיוצרים שאילתה מתמשכת, צריך להשתמש בקידומת מותאמת אישית של מזהה עבודה. כמה שאילתות רציפות יכולות לחלוק את אותו הקידומת. לדוגמה, אפשר להשתמש בקידומת prod- כדי לציין שאילתה היא אילתת ייצור.
  2. נכנסים לדף Log-based Metrics במסוף Cloud de Confiance .

    כניסה אל Log-based Metrics

  3. לוחצים על יצירת מדד. החלונית Create logs metric תופיע.

  4. בקטע סוג המדד, בוחרים באפשרות מונה.

  5. בקטע פרטים, נותנים שם למדד. לדוגמה, CUSTOM_JOB_ID_PREFIX-metric.

  6. בקטע Filter selection, מזינים את הטקסט הבא בכלי לעריכת Build filter:

    resource.type = "bigquery_project"
    protoPayload.resourceName : "projects/PROJECT_ID/jobs/CUSTOM_JOB_ID_PREFIX"
    severity = ERROR
    

    מחליפים את מה שכתוב בשדות הבאים:

  7. לוחצים על יצירת מדד.

  8. בתפריט הניווט, לוחצים על מדדים מבוססי-יומן. המדד שיצרתם מופיע ברשימת המדדים שהוגדרו על ידי המשתמש.

  9. בשורה של המדד, לוחצים על עוד פעולות ואז על יצירת התראה מהמדד.

  10. לוחצים על הבא. אין צורך לשנות את הגדרות ברירת המחדל בדף Policy configuration mode (מצב הגדרת מדיניות).

  11. לוחצים על הבא. אין צורך לשנות את הגדרות ברירת המחדל בדף Configure alert trigger (הגדרת טריגר להתראה).

  12. בוחרים את ערוצי ההתראות ומזינים שם למדיניות ההתראות.

  13. לוחצים על יצירת מדיניות.

כדי לבדוק את ההתראה, מריצים שאילתה מתמשכת עם הקידומת המותאמת אישית של מזהה העבודה שבחרתם, ואז מבטלים אותה. יכול להיות שיחלפו כמה דקות עד שההתראה תגיע לערוץ ההתראות.

ניסיון חוזר של שאילתות שנכשלו

ניסיון חוזר להפעיל שאילתה מתמשכת שנכשלה יכול לעזור למנוע מצבים שבהם פייפליין רציף מושבת למשך זמן ממושך או שנדרשת התערבות אנושית כדי להפעיל אותו מחדש. כשמנסים להריץ מחדש שאילתה מתמשכת שנכשלה, חשוב לקחת בחשבון את הדברים הבאים:

  • האם אפשר לעבד מחדש חלק מהנתונים שעובדו על ידי השאילתה הקודמת לפני שהיא נכשלה.
  • איך לטפל בהגבלת ניסיונות חוזרים או להשתמש בהשהיה מעריכית לפני ניסיון חוזר (exponential backoff).

אחת מהגישות האפשריות לאוטומציה של ניסיון חוזר של שאילתה היא:

  1. יוצרים יעד ב-Cloud Logging על סמך מסנן הכללה שתואם לקריטריונים הבאים, כדי לנתב יומנים לנושא Pub/Sub:

    resource.type = "bigquery_project"
    protoPayload.resourceName : "projects/PROJECT_ID/jobs/CUSTOM_JOB_ID_PREFIX"
    severity = ERROR
    

    מחליפים את מה שכתוב בשדות הבאים:

  2. יוצרים פונקציית Cloud Run שמופעלת בתגובה ליומנים שתואמים למסנן שלכם שמתקבלים ב-Pub/Sub.

    פונקציית Cloud Run יכולה לקבל את מטען הנתונים מהודעת Pub/Sub ולנסות להפעיל שאילתה רציפה חדשה באמצעות אותה תחביר SQL כמו השאילתה שנכשלה, אבל בתחילת התהליך, מיד אחרי שהעבודה הקודמת נעצרה.

לדוגמה, אתם יכולים להשתמש בפונקציה שדומה לזו:

Python

לפני שמנסים את הדוגמה הזו, צריך לפעול לפי Pythonהוראות ההגדרה שבמדריך למתחילים של BigQuery באמצעות ספריות לקוח. מידע נוסף מופיע במאמרי העזרה של BigQuery Python API.

כדי לבצע אימות ב-BigQuery, צריך להגדיר את Application Default Credentials. מידע נוסף זמין במאמר הגדרת אימות לספריות לקוח.

לפני שמריצים דוגמאות קוד, צריך להגדיר את משתנה הסביבה GOOGLE_CLOUD_UNIVERSE_DOMAIN לערך s3nsapis.fr.

import base64
import json
import logging
import re
import uuid

import google.auth
import google.auth.transport.requests
import requests


def retry_continuous_query(event, context):
    logging.info("Cloud Function started.")

    if "data" not in event:
        logging.info("No data in Pub/Sub message.")
        return

    try:
        # Decode and parse the Pub/Sub message data
        log_entry = json.loads(base64.b64decode(event["data"]).decode("utf-8"))

        # Extract the SQL query and other necessary data
        proto_payload = log_entry.get("protoPayload", {})
        metadata = proto_payload.get("metadata", {})
        job_change = metadata.get("jobChange", {})
        job = job_change.get("job", {})
        job_config = job.get("jobConfig", {})
        query_config = job_config.get("queryConfig", {})
        sql_query = query_config.get("query")
        job_stats = job.get("jobStats", {})
        end_timestamp = job_stats.get("endTime")
        failed_job_id = job.get("jobName")

        # Check if required fields are missing
        if not all([sql_query, failed_job_id, end_timestamp]):
            logging.error("Required fields missing from log entry.")
            return

        logging.info(f"Retrying failed job: {failed_job_id}")

        # Adjust the timestamp in the SQL query
        timestamp_match = re.search(
            r"\s*TIMESTAMP\(('.*?')\)(\s*\+ INTERVAL 1 MICROSECOND)?", sql_query
        )

        if timestamp_match:
            original_timestamp = timestamp_match.group(1)
            new_timestamp = f"'{end_timestamp}'"
            sql_query = sql_query.replace(original_timestamp, new_timestamp)
        elif "CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE" in sql_query:
            new_timestamp = f"TIMESTAMP('{end_timestamp}') + INTERVAL 1 MICROSECOND"
            sql_query = sql_query.replace(
                "CURRENT_TIMESTAMP() - INTERVAL 10 MINUTE", new_timestamp
            )

        # Get access token
        credentials, project = google.auth.default(
            scopes=["https://www.googleapis.com/auth/cloud-platform"]
        )
        request = google.auth.transport.requests.Request()
        credentials.refresh(request)
        access_token = credentials.token

        # API endpoint
        url = f"https://bigquery.googleapis.com/bigquery/v2/projects/{project}/jobs"

        # Request headers
        headers = {
            "Authorization": f"Bearer {access_token}",
            "Content-Type": "application/json",
        }

        # Generate a random UUID
        random_suffix = str(uuid.uuid4())[:8]  # Take the first 8 characters of the UUID

        # Combine the prefix and random suffix
        job_id = f"CUSTOM_JOB_ID_PREFIX{random_suffix}"

        # Request payload
        data = {
            "configuration": {
                "query": {
                    "query": sql_query,
                    "useLegacySql": False,
                    "continuous": True,
                    "connectionProperties": [
                        {"key": "service_account", "value": "SERVICE_ACCOUNT"}
                    ],
                    # ... other query parameters ...
                },
                "labels": {"bqux_job_id_prefix": "CUSTOM_JOB_ID_PREFIX"},
            },
            "jobReference": {
                "projectId": project,
                "jobId": job_id,  # Use the generated job ID here
            },
        }

        # Make the API request
        response = requests.post(url, headers=headers, json=data)

        # Handle the response
        if response.status_code == 200:
            logging.info("Query job successfully created.")
        else:
            logging.error(f"Error creating query job: {response.text}")

    except Exception as e:
        logging.error(
            f"Error processing log entry or retrying query: {e}", exc_info=True
        )

    logging.info("Cloud Function finished.")

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