ייבוא, ייצוא ושינוי של נתונים באמצעות Dataflow

בדף הזה מוסבר איך להשתמש במחבר Dataflow ל-Spanner כדי לייבא, לייצא ולשנות נתונים במסדי נתונים של Spanner ב-GoogleSQL וב-PostgreSQL.

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

מחבר Dataflow הוא השיטה המומלצת להעברת נתונים אל Spanner וממנו בכמות גדולה בצורה יעילה. זו גם השיטה המומלצת לביצוע טרנספורמציות גדולות במסד נתונים שלא נתמכות על ידי Partitioned DML, כמו העברות של טבלאות ומחיקות בכמות גדולה שדורשות JOIN. כשעובדים עם מסדי נתונים בודדים, אפשר להשתמש בשיטות אחרות כדי לייבא ולייצא נתונים:

  • אפשר להשתמש במסוף Google Cloud כדי לייצא מסד נתונים בודד מ-Spanner אל Cloud Storage בפורמט Avro.
  • אפשר להשתמש במסוף Google Cloud כדי לייבא מסד נתונים בחזרה ל-Spanner מקבצים שייצאתם ל-Cloud Storage.
  • אפשר להשתמש ב-API בארכיטקטורת REST או ב-Google Cloud CLI כדי להריץ משימות ייצוא או ייבוא מ-Spanner ל-Cloud Storage ובחזרה, גם באמצעות פורמט Avro.

המחבר של Dataflow ל-Spanner הוא חלק מ-Apache Beam Java SDK, והוא מספק API לביצוע הפעולות הקודמות. למידע נוסף על חלק מהמושגים שמוסברים בדף הזה, כמו אובייקטים וטרנספורמציות, אפשר לעיין במדריך לתכנות ב-Apache Beam.PCollection

הוספת המחבר לפרויקט Maven

כדי להוסיף את מחבר Dataflow לפרויקט Maven, מוסיפים את ארטיפקט Maven לקובץ pom.xml כתלות. Google Cloud beam-sdks-java-io-google-cloud-platform

לדוגמה, אם בקובץ pom.xml מוגדר beam.version למספר הגרסה המתאים, מוסיפים את התלות הבאה:

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-google-cloud-platform</artifactId>
    <version>${beam.version}</version>
</dependency>

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

כדי לקרוא מ-Spanner, צריך להחיל את טרנספורמציית SpannerIO.read. מגדירים את הקריאה באמצעות השיטות בכיתה SpannerIO.Read. הפעלת הטרנספורמציה מחזירה PCollection<Struct>, כאשר כל רכיב באוסף מייצג שורה נפרדת שהוחזרה על ידי פעולת הקריאה. אתם יכולים לקרוא מ-Spanner עם שאילתת SQL ספציפית או בלי שאילתת SQL ספציפית, בהתאם לפלט שאתם צריכים.

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

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

כדי לקרוא קבוצה ספציפית של נתונים מ-Spanner, צריך להגדיר את הטרנספורמציה באמצעות השיטה SpannerIO.Read.withQuery כדי לציין שאילתת SQL. לדוגמה:

// Query for all the columns and rows in the specified Spanner table
PCollection<Struct> records = pipeline.apply(
    SpannerIO.read()
        .withInstanceId(instanceId)
        .withDatabaseId(databaseId)
        .withQuery("SELECT * FROM " + options.getTable()));

קריאת נתונים בלי לציין שאילתה

כדי לקרוא מתוך מסד נתונים בלי להשתמש בשאילתה, אפשר לציין שם של טבלה באמצעות ה-method ‏SpannerIO.Read.withTable, ולציין רשימה של עמודות לקריאה באמצעות ה-method ‏SpannerIO.Read.withColumns. לדוגמה:

GoogleSQL

// Query for all the columns and rows in the specified Spanner table
PCollection<Struct> records = pipeline.apply(
    SpannerIO.read()
        .withInstanceId(instanceId)
        .withDatabaseId(databaseId)
        .withTable("Singers")
        .withColumns("singerId", "firstName", "lastName"));

PostgreSQL

// Query for all the columns and rows in the specified Spanner table
PCollection<Struct> records = pipeline.apply(
    SpannerIO.read()
        .withInstanceId(instanceId)
        .withDatabaseId(databaseId)
        .withTable("singers")
        .withColumns("singer_id", "first_name", "last_name"));

כדי להגביל את מספר השורות שנקראות, אפשר לציין קבוצה של מפתחות ראשיים לקריאה באמצעות ה-method‏ SpannerIO.Read.withKeySet.

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

כדי לעשות זאת, מציינים את הטבלה כמו בדוגמה הקודמת, ומציינים את האינדקס שמכיל את ערכי העמודות הנדרשים באמצעות השיטה SpannerIO.Read.withIndex. האינדקס צריך לאחסן את כל העמודות שהטרנספורמציה צריכה לקרוא. המפתח הראשי של טבלת הבסיס מאוחסן באופן מרומז. לדוגמה, כדי לקרוא את הטבלה Songs באמצעות האינדקס SongsBySongName, משתמשים בקוד הבא:

GoogleSQL

// Read the indexed columns from all rows in the specified index.
PCollection<Struct> records =
    pipeline.apply(
        SpannerIO.read()
            .withInstanceId(instanceId)
            .withDatabaseId(databaseId)
            .withTable("Songs")
            .withIndex("SongsBySongName")
            // Can only read columns that are either indexed, STORED in the index or
            // part of the primary key of the Songs table,
            .withColumns("SingerId", "AlbumId", "TrackId", "SongName"));

PostgreSQL

// // Read the indexed columns from all rows in the specified index.
PCollection<Struct> records =
    pipeline.apply(
        SpannerIO.read()
            .withInstanceId(instanceId)
            .withDatabaseId(databaseId)
            .withTable("Songs")
            .withIndex("SongsBySongName")
            // Can only read columns that are either indexed, STORED in the index or
            // part of the primary key of the songs table,
            .withColumns("singer_id",