הצגת שרשרת מקורות נתונים במערכות Google Cloud

אפשר להציג את שושלת הנתונים כדי להבין את הקשרים בין המשאבים של הפרויקט לבין התהליכים שיצרו אותם. הקשרים האלה מראים איך נכסי נתונים, כמו טבלאות ומערכי נתונים, עוברים טרנספורמציה בתהליכים כמו שאילתות וצינורות נתונים. במדריך הזה מוסבר איך לראות את הפרטים של שרשרת מקורות הנתונים במסוף Google Cloud או לאחזר אותם באמצעות Data Lineage API.

תפקידים והרשאות

כשמפעילים את Data Lineage API, המערכת עוקבת אחרי מידע על מקור הנתונים באופן אוטומטי. לא צריך הרשאות אדמין או עריכה כדי לתעד את מקורות הנתונים של נכסי הנתונים.

כדי לראות את שרשרת מקורות הנתונים, אתם צריכים הרשאות ספציפיות לניהול זהויות והרשאות גישה (IAM). פרטי השושלת נאספים בכל הפרויקטים, ולכן צריך הרשאות בכמה פרויקטים.

  • כשמציגים את היסטוריית השינויים ב-Knowledge Catalog, ב-BigQuery או ב-Vertex AI: צריך הרשאות להצגת פרטים על היסטוריית השינויים בפרויקט שבו מציגים אותה.

  • כשמציגים את היסטוריית השינויים שנרשמה בפרויקטים אחרים: צריך הרשאות להצגת פרטי היסטוריית השינויים בפרויקטים שבהם היא נרשמה.

כדי לקבל את ההרשאות שדרושות בשביל להציג את היסטוריית הנתונים, אתם צריכים לבקש מהאדמין לתת לכם את התפקידים הבאים ב-IAM:

  • Data Lineage Viewer (roles/datalineage.viewer) בפרויקט שבו מתועד שרשרת המקורות, ובפרויקט שבו מוצגת שרשרת המקורות
  • צפייה בפרטי הטבלה ב-BigQuery: ‫BigQuery Data Viewer (roles/bigquery.dataViewer) בפרויקט האחסון של הטבלה
  • צפייה בפרטי משימה ב-BigQuery: BigQuery Resource Viewer (roles/bigquery.resourceViewer) בפרויקט המחשוב של המשימה
  • הצגת פרטים של נכסים אחרים בקטלוג: Dataplex Catalog Viewer (roles/dataplex.catalogViewer) בפרויקט שבו מאוחסנים רשומות הקטלוג

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

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

ההרשאות הנדרשות

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

  • הצגת פרטים של טבלה ב-BigQuery: ‫bigquery.tables.get – פרויקט האחסון של הטבלה
  • הצגת פרטי המשימה ב-BigQuery: bigquery.jobs.get – פרויקט המחשוב של המשימה

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

סוגים של תצוגות של שרשרת מקורות הנתונים

אפשר לראות את פרטי השושלת בתרשים אינטראקטיבי או ברשימה מובנית במסוף Google Cloud .

תיאור מפורט של רכיבי התרשים (כמו צמתים, קשתות, סמלי תהליך ותוויות) והעמודות שזמינות בתצוגות הרשימה מופיע במאמר מידע על ויזואליזציה של השתלשלות הנתונים ב-Knowledge Catalog.

הפעלת שושלת נתונים

מפעילים את תכונת שרשרת המקורות כדי להתחיל לעקוב באופן אוטומטי אחרי מידע על שרשרת המקורות במערכות נתמכות. כברירת מחדל, הפעלת ה-API מפעילה מעקב אחר מקורות נתונים ברוב השירותים הנתמכים. כדי לשלוט בהטמעת שושלת הנתונים ב-Managed Service for Apache Spark, אפשר לעיין במאמר בנושא שליטה בהטמעת שושלת נתונים בשירות.

צריך להפעיל את Data Lineage API גם בפרויקט שבו צופים ב-lineage וגם בפרויקטים שבהם מתועד ה-lineage. מידע נוסף זמין במאמר בנושא סוגי פרויקטים.

  1. כדי לתעד את פרטי השושלת:
    1. בדף Project selector במסוף Google Cloud , בוחרים את הפרויקט שבו רוצים לתעד את היחסים בין מקורות הנתונים.

      כניסה לדף לבחירת הפרויקט

    2. מפעילים את Data Lineage API.

      הפעלת ה-API

    3. חוזרים על השלבים הקודמים לכל פרויקט שבו רוצים לתעד את היסטוריית השינויים.
  2. בפרויקט שבו צופים ב-Data Lineage, מפעילים את Data Lineage API ואת Dataplex API.

    הפעלת ממשקי API

שליטה בהטמעת שרשרת המקורות בשביל שירות

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

לפרטים על אופן היישום ההיררכי של ההגדרות האלה דרך עץ המשאבים, אפשר לעיין במאמר שליטה בהטמעת שושלת הנתונים.

הצגת שושלת נתונים

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

המסוף

אפשר לגשת למידע על מקורות הנתונים במסוף Google Cloud מנקודות התחלה שונות:

  • Knowledge Catalog: עוברים לדף חיפוש של Knowledge Catalog, בוחרים באפשרות Knowledge Catalog כמצב החיפוש, מחפשים את הרשומה שרוצים לראות ואז לוחצים עליה. מידע נוסף זמין במאמר חיפוש משאבים ב-Knowledge Catalog.
  • BigQuery: עוברים לדף BigQuery ופותחים את הטבלה שרוצים לראות את שושלת הנתונים שלה.
  • Vertex AI: עוברים לדף Datasets או Model Registry ולוחצים על מערך הנתונים או על המודל שרוצים לראות את שרשרת המקורות שלהם.

כדי לראות את גרף שושלת הנתונים, בצע את השלבים הבאים:

  1. לוחצים על הכרטיסייה Lineage (מקורות נתונים).

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

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

    מידע נוסף זמין במאמר בנושא עיון ידני בתרשים שושלת הנתונים.

  3. לוחצים על צומת בתצוגה Graph.

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

  4. לוחצים על קצה עם סמל של תהליך בתצוגה תרשים.

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

    • כדי לבדוק את לוגיקת השינוי, לוחצים על הכרטיסייה פרטים.
    • כדי לראות את הביקורת ואת היסטוריית ההרצות, לוחצים על הכרטיסייה Runs (הרצות).
  5. בחלונית Lineage explorer, בוחרים קריטריונים לסינון – לדוגמה, Direction,‏ Dependency type או Time range – ואז לוחצים על Apply.

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

  6. בתצוגה הממוקדת Graph, בוחרים צומת, ואז בחלונית הפרטים של הצומת לוחצים על Visualize Path כדי להציג את נתיב השושלת מהצומת שנבחר בחזרה אל רשומת הבסיס (רק בתצוגה הממוקדת).

    מידע נוסף זמין במאמר בנושא הדמיה של נתיב שושלת.

  7. כדי לראות את שושלת הנתונים ברמת העמודה (רק למשימות של BigQuery ו-Managed Service for Apache Spark), מבצעים אחת מהפעולות הבאות:

    • בתצוגת תרשים ממוקדת, לוחצים על סמל העמודה בטבלה.
      הסמל שמשמש למעבר לנתוני שושלת ברמת העמודה.
      סמל העמודות
    • בחלונית Lineage explorer (כלי לבדיקת מקורות נתונים), מסננים לפי שם העמודה ולוחצים על Apply (החלה).

    מידע נוסף זמין במאמר בנושא Column-level lineage (היסטוריה ברמת העמודה).

  8. לוחצים על איפוס.

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

  9. לוחצים על רשימה כדי לעבור לתצוגת הרשימה.

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

Java

import com.google.api.gax.rpc.ApiException;
import com.google.cloud.datacatalog.lineage.v1.BatchSearchLinkProcessesRequest;
import com.google.cloud.datacatalog.lineage.v1.EntityReference;
import com.google.cloud.datacatalog.lineage.v1.EventLink;
import com.google.cloud.datacatalog.lineage.v1.LineageClient;
import com.google.cloud.datacatalog.lineage.v1.LineageEvent;
import com.google.cloud.datacatalog.lineage.v1.Link;
import com.google.cloud.datacatalog.lineage.v1.ListLineageEventsRequest;
import com.google.cloud.datacatalog.lineage.v1.ListRunsRequest;
import com.google.cloud.datacatalog.lineage.v1.LocationName;
import com.google.cloud.datacatalog.lineage.v1.ProcessLinks;
import com.google.cloud.datacatalog.lineage.v1.Run;
import com.google.cloud.datacatalog.lineage.v1.SearchLinksRequest;
import java.io.IOException;
import java.util.ArrayList;
import java.util.HashSet;
import java.util.LinkedList;
import java.util.List;
import java.util.Queue;
import java.util.Set;

public class ViewLineageExample {

  public static void main(String[] args) throws IOException {
    // TODO(developer): Replace these variables before running the sample.
    String projectId = "my-project-id";
    String location = "us";
    String targetFullyQualifiedName = "bigquery:my-project-id.my_dataset.my_table";
    int maxDepth = 3;

    viewLineage(projectId, location, targetFullyQualifiedName, maxDepth);
  }

  static class Node {
    String fqn;
    int depth;
    Node(String fqn, int depth) {
      this.fqn = fqn;
      this.depth = depth;
    }
  }

  public static void viewLineage(
      String projectId, String location, String targetFullyQualifiedName, int maxDepth)
      throws IOException {
    // Initialize client that will be used to send requests. This client only needs
    // to be created once, and can be reused for multiple requests.
    try (LineageClient client = LineageClient.create()) {
      String parent = LocationName.of(projectId, location).toString();

      Set<String> visitedNodes = new HashSet<>();
      Queue<Node> queue = new LinkedList<>();

      visitedNodes.add(targetFullyQualifiedName);
      queue.offer(new Node(targetFullyQualifiedName, 0));

      while (!queue.isEmpty()) {
        Node current = queue.poll();
        System.out.printf("\nExploring node (Depth %d): %s\n", current.depth, current.fqn);

        if (current.depth >= maxDepth) {
          continue;
        }

        EntityReference targetEntity =
            EntityReference.newBuilder().setFullyQualifiedName(current.fqn).build();
        SearchLinksRequest searchLinksRequest =
            SearchLinksRequest.newBuilder().setParent(parent).setTarget(targetEntity).build();

        List<String> linkNames = new ArrayList<>();
        try {
          // 1. Search for links related to the target entity
          for (Link link : client.searchLinks(searchLinksRequest).iterateAll()) {
            linkNames.add(link.getName());
          }
        } catch (ApiException e) {
          System.out.printf("  Failed to retrieve links for %s: %s\n", current.fqn, e.getMessage());
          continue;
        }

        if (linkNames.isEmpty()) {
          continue;
        }

        // 2. Batch search for processes in chunks of 100
        for (int i = 0; i < linkNames.size(); i += 100) {
          List<String> batch = linkNames.subList(i, Math.min(linkNames.size(), i + 100));
          BatchSearchLinkProcessesRequest batchSearchRequest =
              BatchSearchLinkProcessesRequest.newBuilder()
                  .setParent(parent)
                  .addAllLinks(batch)
                  .build();

          try {
            for (ProcessLinks processLinks :
                client.batchSearchLinkProcesses(batchSearchRequest).iterateAll()) {
              String processName = processLinks.getProcess();
              System.out.printf("  Process: %s\n", processName);

              // 3. List runs for the process
              ListRunsRequest runsRequest =
                  ListRunsRequest.newBuilder().setParent(processName).build();
              for (Run run : client.listRuns(runsRequest).iterateAll()) {
                System.out.printf("    Run: %s\n", run.getName());

                // 4. List events for the run
                ListLineageEventsRequest eventsRequest =
                    ListLineageEventsRequest.newBuilder().setParent(run.getName()).build();
                for (LineageEvent event : client.listLineageEvents(eventsRequest).iterateAll()) {
                  for (EventLink eventLink : event.getLinksList()) {
                    String sourceFqn = eventLink.getSource().getFullyQualifiedName();
                    // If exploring upstream, queue the source
                    if (!sourceFqn.isEmpty() && !visitedNodes.contains(sourceFqn)) {
                      visitedNodes.add(sourceFqn);
                      queue.offer(new Node(sourceFqn, current.depth + 1));
                    }
                  }
                }
              }
            }
          } catch (ApiException e) {
            System.out.printf("  Failed to retrieve processes/runs: %s\n", e.getMessage());
          }
        }
      }
    }
  }
}

Python

from google.cloud import datacatalog_lineage_v1
from google.api_core.exceptions import GoogleAPICallError

def view_lineage(project_id: str, location: str, target_fully_qualified_name: str, max_depth: int = 3):
    """Retrieves lineage for a given entity using a depth-limited search."""
    client = datacatalog_lineage_v1.LineageClient()
    parent = f"projects/{project_id}/locations/{location}"

    # Store visited nodes to avoid infinite loops in cyclic graphs
    visited_nodes = set([target_fully_qualified_name])
    queue = [(target_fully_qualified_name, 0)]

    while queue:
        current_node, current_depth = queue.pop(0)
        print(f"\nExploring node (Depth {current_depth}): {current_node}")

        if current_depth >= max_depth:
            continue

        target_entity = datacatalog_lineage_v1.EntityReference(
            fully_qualified_name=current_node
        )
        search_links_request = datacatalog_lineage_v1.SearchLinksRequest(
            parent=parent,
            target=target_entity,
        )

        try:
            links = list(client.search_links(request=search_links_request))
        except GoogleAPICallError as e:
            print(f"  Failed to retrieve links for {current_node}: {e.message}")
            continue

        if not links:
            continue

        # Extract link names to query processes in batches
        link_names = [link.name for link in links]

        # Batch max size is 100
        for i in range(0, len(link_names), 100):
            batch = link_names[i:i + 100]
            batch_request = datacatalog_lineage_v1.BatchSearchLinkProcessesRequest(
                parent=parent,
                links=batch
            )

            try:
                for process_links in client.batch_search_link_processes(request=batch_request):
                    process_name = process_links.process
                    print(f"  Process: {process_name}")

                    runs_request = datacatalog_lineage_v1.ListRunsRequest(parent=process_name)
                    for run in client.list_runs(request=runs_request):
                        print(f"    Run: {run.name}")

                        events_request = datacatalog_lineage_v1.ListLineageEventsRequest(parent=run.name)
                        for event in client.list_lineage_events(request=events_request):
                            for event_link in event.links:
                                source_fqn = event_link.source.fully_qualified_name

                                # If exploring upstream, queue the source
                                if source_fqn and source_fqn not in visited_nodes:
                                    visited_nodes.add(source_fqn)
                                    queue.append((source_fqn, current_depth + 1))

            except GoogleAPICallError as e:
                 print(f"  Failed to retrieve processes/runs: {e.message}")

שיפור ההצגה החזותית של שרשרת המקור

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

  1. כדי לחפש פרויקטים, מערכי נתונים או שמות ישויות ספציפיים, משתמשים בחלונית Filters (מסננים).

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

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

  3. בוחרים אחת מהאפשרויות הבאות או את שתיהן:

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

אפשר לבחור את שתי האפשרויות בו-זמנית. אם שתי האפשרויות מסומנות, הצמתים שלא מסוננים מוסתרים, והצמתים התואמים מודגשים בתצוגת הגרף המסוננת.

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