שימוש ב-Storage Write API ‏ (REST)

במאמר הזה מוסבר איך להזרים נתונים ל-BigQuery באמצעות BigQuery Storage Write API (REST), שנקרא בעבר השיטה הקודמת tabledata.insertAll.

בפרויקטים חדשים מומלץ להשתמש ב-BigQuery Storage Write API‏ (gRPC) במקום ב-Storage Write API‏ (REST). המחיר של Storage Write API ‏ (gRPC) נמוך יותר והוא כולל תכונות חזקות יותר, כולל סמנטיקה של מסירה חד-פעמית וסטרימינג לטבלאות מנוהלות של Apache Iceberg. אם אתם מעבירים פרויקט קיים מ-Storage Write API ‏ (REST) ל-Storage Write API ‏ (gRPC), מומלץ לבחור בזרם ברירת המחדל. ‫Storage Write API (REST) עדיין נתמך באופן מלא.

לפני שמתחילים

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

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

  3. מוודאים שהחיוב מופעל בפרויקט Google Cloud .

  4. הסטרימינג לא זמין דרך התוכנית בחינם. אם תנסו להשתמש בסטרימינג בלי להפעיל את החיוב, תקבלו את השגיאה הבאה: BigQuery: Streaming insert is not allowed in the free tier.

  5. להקצות תפקידים של ניהול זהויות והרשאות גישה (IAM) שנותנים למשתמשים את ההרשאות הנדרשות לביצוע כל משימה שמופיעה במאמר הזה.

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

כדי להזרים נתונים ל-BigQuery, אתם צריכים את הרשאות ה-IAM הבאות:

  • bigquery.tables.updateData (מאפשר להוסיף נתונים לטבלה)
  • ‫bigquery.tables.get (מאפשר לקבל מטא-נתונים של טבלה)
  • ‫bigquery.datasets.get (מאפשר לקבל מטא-נתונים של מערך נתונים)
  • ‫bigquery.tables.create (חובה אם משתמשים בטבלת תבנית כדי ליצור את הטבלה באופן אוטומטי)

כל אחד מהתפקידים הבאים שמוגדרים מראש ב-IAM כולל את ההרשאות שצריך כדי להזרים נתונים ל-BigQuery:

  • roles/bigquery.dataEditor
  • roles/bigquery.dataOwner
  • roles/bigquery.admin

במאמר תפקידים והרשאות מוגדרים מראש יש מידע נוסף על תפקידים והרשאות ב-IAM ב-BigQuery.

הזרמת נתונים ל-BigQuery

C#

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

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


using Google.Cloud.BigQuery.V2;

public class BigQueryTableInsertRows
{
    public void TableInsertRows(
        string projectId = "your-project-id",
        string datasetId = "your_dataset_id",
        string tableId = "your_table_id"
    )
    {
        BigQueryClient client = BigQueryClient.Create(projectId);
        BigQueryInsertRow[] rows = new BigQueryInsertRow[]
        {
            // The insert ID is optional, but can avoid duplicate data
            // when retrying inserts.
            new BigQueryInsertRow(insertId: "row1") {
                { "name", "Washington" },
                { "post_abbr", "WA" }
            },
            new BigQueryInsertRow(insertId: "row2") {
                { "name", "Colorado" },
                { "post_abbr", "CO" }
            }
        };
        client.InsertRows(datasetId, tableId, rows);
    }
}

Go

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

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

import (
	"context"
	"fmt"

	"cloud.google.com/go/bigquery"
)

// Item represents a row item.
type Item struct {
	Name string
	Age  int
}

// Save implements the ValueSaver interface.
// This example disables best-effort de-duplication, which allows for higher throughput.
func (i *Item) Save() (map[string]bigquery.Value, string, error) {
	return map[string]bigquery.Value{
		"full_name": i.Name,
		"age":       i.Age,
	}, bigquery.NoDedupeID, nil
}

// insertRows demonstrates inserting data into a table using the streaming insert mechanism.
func insertRows(projectID, datasetID, tableID string) error {
	// projectID := "my-project-id"
	// datasetID := "mydataset"
	// tableID := "mytable"
	ctx := context.Background()
	client, err := bigquery.NewClient(ctx, projectID)
	if err != nil {
		return fmt.Errorf("bigquery.NewClient: %w", err)
	}
	defer client.Close()

	inserter := client.Dataset(datasetID).Table(tableID).Inserter()
	items := []*Item{
		// Item implements the ValueSaver interface.
		{Name: "Phred Phlyntstone", Age: 32},
		{Name: "Wylma Phlyntstone", Age: 29},
	}
	if err := inserter.Put(ctx, items); err != nil {
		return err
	}
	return nil
}

Java

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

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

import com.google.cloud.bigquery.BigQuery;
import com.google.cloud.bigquery.BigQueryError;
import com.google.cloud.bigquery.BigQueryException;
import com.google.cloud.bigquery.BigQueryOptions;
import com.google.cloud.bigquery.InsertAllRequest;
import com.google.cloud.bigquery.InsertAllResponse;
import com.google.cloud.bigquery.TableId;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

// Sample to inserting rows into a table without running a load job.
public class TableInsertRows {

  public static void main(String[] args) {
    // TODO(developer): Replace these variables before running the sample.
    String datasetName = "MY_DATASET_NAME";
    String tableName = "MY_TABLE_NAME";
    // Create a row to insert
    Map, Object> rowContent = new HashMap<>();
    rowContent.put("booleanField", true);
    rowContent.put("numericField", "3.14");
    // TODO(developer): Replace the row id with a unique value for each row.
    String rowId = "ROW_ID";
    tableInsertRows(datasetName, tableName, rowId, rowContent);
  }

  public static void tableInsertRows(
      String datasetName, String tableName, String rowId, Map, Object> rowContent) {
    try {
      // Initialize client that will be used to send requests. This client only needs to be created
      // once, and can be reused for multiple requests.
      BigQuery bigquery = BigQueryOptions.getDefaultInstance().getService();

      // Get table
      TableId tableId = TableId.of(datasetName, tableName);

      // Inserts rowContent into datasetName:tableId.
      InsertAllResponse response =
          bigquery.insertAll(
              InsertAllRequest.newBuilder(tableId)
                  // More rows can be added in the same RPC by invoking .addRow() on the builder.
                  // You can omit the unique row ids to disable de-duplication.
                  .addRow(rowId, rowContent)
                  .build());

      if (response.hasErrors()) {
        // If any of the insertions failed, this lets you inspect the errors
        for (Map.Entry, List> entry : response.getInsertErrors().entrySet()) {
          System.out.println("Response error: \n" + entry.getValue());
        }
      }
      System.out.println("Rows successfully inserted into table");
    } catch (BigQueryException e) {
      System.out.println("Insert operation not performed \n" + e.toString());
    }
  }
}

Node.js

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

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

// Import the Google Cloud client library
const {BigQuery} = require('@google-cloud/bigquery');
const bigquery = new BigQuery();

async function insertRowsAsStream() {
  // Inserts the JSON objects into my_dataset:my_table.

  /**
   * TODO(developer): Uncomment the following lines before running the sample.
   */
  // const datasetId = 'my_dataset';
  // const tableId = 'my_table';
  const rows = [
    {name: 'Tom', age: 30},
    {name: 'Jane', age: 32},
  ];

  // Insert data into a table
  await bigquery.dataset(datasetId).table(tableId).insert(rows);
  console.log(`Inserted ${rows.length} rows`);
}

PHP

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

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

use Google\Cloud\BigQuery\BigQueryClient;

/**
 * Stream data into bigquery
 *
 * @param string $projectId The project Id of your Google Cloud Project.
 * @param string $datasetId The BigQuery dataset ID.
 * @param string $tableId The BigQuery table ID.
 * @param string $data Json encoded data For eg,
 *    $data = json_encode([
 *       "field1" => "value1",
 *       "field2" => "value2",
 *    ]);
 */
function stream_row(
    string $projectId,
    string $datasetId,
    string $tableId,
    string $data
): void {
    // instantiate the bigquery table service
    $bigQuery = new BigQueryClient([
      'projectId' => $projectId,
    ]);
    $dataset = $bigQuery->dataset($datasetId);
    $table = $dataset->table($tableId);

    $data = json_decode($data, true);
    $insertResponse = $table->insertRows([
      ['data' => $data],
      // additional rows can go here
    ]);
    if ($insertResponse->isSuccessful()) {
        print('Data streamed into BigQuery successfully' . PHP_EOL);
    } else {
        foreach ($insertResponse->failedRows() as $row) {
            foreach ($row['errors'] as $error) {
                printf('%s: %s' . PHP_EOL, $error['reason'], $error['message']);
            }
        }
    }
}

Python

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

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

from google.cloud import bigquery

# Construct a BigQuery client object.
client = bigquery.Client()

# TODO(developer): Set table_id to the ID of table to append to.
# table_id = "your-project.your_dataset.your_table"

rows_to_insert = [
    {"full_name": "Phred Phlyntstone", "age": 32},
    {"full_name": "Wylma Phlyntstone", "age": 29},
]

errors = client.insert_rows_json(table_id, rows_to_insert)  # Make an API request.
if errors == []:
    print("New rows have been added.")
else:
    print("Encountered errors while inserting rows: {}".format(errors))

Ruby

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

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

require "google/cloud/bigquery"

def table_insert_rows dataset_id = "your_dataset_id", table_id = "your_table_id"
  bigquery = Google::Cloud::Bigquery.new
  dataset  = bigquery.dataset dataset_id
  table    = dataset.table table_id

  row_data = [
    { name: "Alice", value: 5  },
    { name: "Bob",   value: 10 }
  ]
  response = table.insert row_data

  if response.success?
    puts "Inserted rows successfully"
  else
    puts "Failed to insert #{response.error_rows.count} rows"
  end
end

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

Java

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

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

import com.google.cloud.bigquery.BigQuery;
import com.google.cloud.bigquery.BigQueryError;
import com.google.cloud.bigquery.BigQueryException;
import com.google.cloud.bigquery.BigQueryOptions;
import com.google.cloud.bigquery.InsertAllRequest;
import com.google.cloud.bigquery.InsertAllResponse;
import com.google.cloud.bigquery.TableId;
import com.google.common.collect.ImmutableList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

// Sample to insert rows without row ids in a table
public class TableInsertRowsWithoutRowIds {

  public static void main(String[] args) {
    // TODO(developer): Replace these variables before running the sample.
    String datasetName = "MY_DATASET_NAME";
    String tableName = "MY_TABLE_NAME";
    tableInsertRowsWithoutRowIds(datasetName, tableName);
  }

  public static void tableInsertRowsWithoutRowIds(String datasetName, String tableName) {
    try {
      // Initialize client that will be used to send requests. This client only needs to be created
      // once, and can be reused for multiple requests.
      final BigQuery bigquery = BigQueryOptions.getDefaultInstance().getService();
      // Create rows to insert
      Map, Object> rowContent1 = new HashMap<>();
      rowContent1.put("stringField", "Phred Phlyntstone");
      rowContent1.put("numericField", 32);
      Map, Object> rowContent2 = new HashMap<>();
      rowContent2.put("stringField", "Wylma Phlyntstone");
      rowContent2.put("numericField", 29);
      InsertAllResponse response =
          bigquery.insertAll(
              InsertAllRequest.newBuilder(TableId.of(datasetName, tableName))
                  // No row ids disable de-duplication, and also disable the retries in the Java
                  // library.
                  .setRows(
                      ImmutableList.of(
                          InsertAllRequest.RowToInsert.of(rowContent1),
                          InsertAllRequest.RowToInsert.of(rowContent2)))
                  .build());

      if (response.hasErrors()) {
        // If any of the insertions failed, this lets you inspect the errors
        for (Map.Entry, List> entry : response.getInsertErrors().entrySet()) {
          System.out.println("Response error: \n" + entry.getValue());
        }
      }
      System.out.println("Rows successfully inserted into table without row ids");
    } catch (BigQueryException e) {
      System.out.println("Insert operation not performed \n" + e.toString());
    }
  }
}

Python

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

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

from google.cloud import bigquery

# Construct a BigQuery client object.
client = bigquery.Client()

# TODO(developer): Set table_id to the ID of table to append to.
# table_id = "your-project.your_dataset.your_table"

rows_to_insert = [
    {"full_name": "Phred Phlyntstone", "age": 32},
    {"full_name": "Wylma Phlyntstone", "age": 29},
]

errors = client.insert_rows_json(
    table_id, rows_to_insert, row_ids=[None] * len(rows_to_insert)
)  # Make an API request.
if errors == []:
    print("New rows have been added.")
else:
    print("Encountered errors while inserting rows: {}".format(errors))

שליחת נתונים של תאריך ושעה

בשדות של תאריך ושעה, הפורמט של הנתונים ב-Storage Write API ‏ (REST) הוא:

סוג פורמט
DATE מחרוזת בפורמט "YYYY-MM-DD"
DATETIME מחרוזת בפורמט "YYYY-MM-DD [HH:MM:SS]"
TIME מחרוזת בפורמט "HH:MM:SS"
TIMESTAMP מספר השניות מאז 1970-01-01 (ראשית זמן יוניקס), או מחרוזת בפורמט "YYYY-MM-DD HH:MM[:SS]"

שליחת נתוני טווח

בשדות עם סוג RANGE, צריך לעצב את הנתונים ב-Storage Write API (REST) כאובייקט JSON עם שני שדות: start ו-end. ערכים חסרים או ערכי NULL בשדות start ו-end מייצגים גבולות לא מוגבלים. השדות האלה צריכים להיות בפורמט JSON הנתמך מסוג T, כאשר T יכול להיות אחד מהערכים DATE, DATETIME ו-TIMESTAMP.

בדוגמה הבאה, השדה f_range_date מייצג עמודה RANGE בטבלה. שורה מוכנסת לעמודה הזו באמצעות Storage Write API (REST).

{
    "f_range_date": {
        "start": "1970-01-02",
        "end": null
    }
}

זמינות נתונים בזמן אמת

הנתונים זמינים לניתוח בזמן אמת באמצעות שאילתות GoogleSQL מיד אחרי ש-BigQuery מאשר בהצלחה בקשה של Storage Write API ‏ (REST). כשמריצים שאילתה על נתונים במאגר הזמני של הנתונים, לא מחויבים על בייטים שעובדו מהמאגר הזמני אם משתמשים בתמחור על פי דרישה של משאבי מחשוב. אם אתם משתמשים בתמחור לפי קיבולת, ההזמנות שלכם צורכות משבצות לעיבוד נתונים במאגר הזמני של הנתונים בזמן שהם מועברים.

לשורות שנוספו לאחרונה לטבלה עם חלוקה למחיצות לפי זמן ההוספה יש באופן זמני ערך NULL בעמודה הווירטואלית _PARTITIONTIME. בשביל שורות כאלה, BigQuery מקצה את הערך הסופי שאינו NULL של העמודה PARTITIONTIME ברקע, בדרך כלל תוך כמה דקות. במקרים נדירים, התהליך הזה יכול להימשך עד 90 דקות.

יכול להיות שחלק מהשורות שהוזרמו לאחרונה לא יהיו זמינות להעתקת הטבלה, בדרך כלל למשך כמה דקות. במקרים נדירים, התהליך הזה יכול להימשך עד 90 דקות. כדי לראות אם הנתונים זמינים להעתקת הטבלה, בודקים את התגובה tables.get בקטע שנקרא streamingBuffer. אם הקטע streamingBuffer לא מופיע, הנתונים שלכם זמינים להעתקה. אפשר גם להשתמש בשדה streamingBuffer.oldestEntryTime כדי לזהות את גיל הרשומות במאגר הזמני של הסטרימינג.

הסרת כפילויות בהקדם האפשרי

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

המערכת מצפה שהשורות שמופיעות עם ערכי insertId זהים יהיו גם זהות. אם לשתי שורות יש ערכים זהים ב-insertId, לא ניתן לקבוע איזו שורה BigQuery ישמור.

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

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

אם יש לכם דרישות מחמירות לביטול כפילויות בנתונים, Google Cloud Datastore הוא שירות חלופי שתומך בטרנזקציות.

השבתה של ביטול כפילויות בהקדם האפשרי

כדי להשבית את הסרת הכפילויות בשיטת הכי טוב שאפשר, אל תמלאו את השדה insertId בכל שורה שמוסיפים. זו הדרך המומלצת להוספת נתונים.

‫Apache Beam ו-Dataflow

כדי להשבית את ביטול הכפילויות בשיטת הכי טוב שאפשר כשמשתמשים במחבר BigQuery I/O של Apache Beam ל-Java, צריך להשתמש ב-method‏ ignoreInsertIds().

הסרה ידנית של כפילויות

כדי לוודא שלא יהיו שורות כפולות אחרי סיום הסטרימינג, אפשר לבצע את התהליך הידני הבא:

  1. מוסיפים את insertId כעמודה בסכימת הטבלה וכוללים את הערך insertId בנתונים של כל שורה.
  2. אחרי שהסטרימינג מפסיק, מריצים את השאילתה הבאה כדי לבדוק אם יש כפילויות:
    #standardSQL
    SELECT
      MAX(count) FROM(
      SELECT
        ID_COLUMN,
        count(*) as count
      FROM
        `TABLE_NAME`
      GROUP BY
        ID_COLUMN)
    אם התוצאה גדולה מ-1, יש כפילויות.
  3. כדי להסיר כפילויות, מריצים את השאילתה הבאה. מציינים טבלת יעד, מאפשרים תוצאות גדולות ומשביתים את השטחת התוצאות.
    #standardSQL
    SELECT
      * EXCEPT(row_number)
    FROM (
      SELECT
        *,
        ROW_NUMBER()
              OVER (PARTITION BY ID_COLUMN) row_number
      FROM
        `TABLE_NAME`)
    WHERE
      row_number = 1

הערות לגבי השאילתה להסרת תוכן כפול:

  • האסטרטגיה הבטוחה יותר לשאילתת הסרת הכפילויות היא טירגוט של טבלה חדשה. לחלופין, אפשר לטרגט את טבלת המקור באמצעות מאפיין של פעולת כתיבה WRITE_TRUNCATE.
  • שאילתת ההסרה של הכפילויות מוסיפה עמודה row_number עם הערך 1 לסוף סכימת הטבלה. השאילתה משתמשת בהצהרה SELECT * EXCEPT מ-GoogleSQL כדי להחריג את העמודה row_number מטבלת היעד. הקידומת #standardSQL מפעילה את GoogleSQL לשאילתה הזו. אפשר גם לבחור שמות עמודות ספציפיים כדי להשמיט את העמודה הזו.
  • כדי לשלוח שאילתות לנתונים פעילים בלי כפילויות, אפשר גם ליצור תצוגה מעל הטבלה באמצעות השאילתה להסרת כפילויות. חשוב לדעת שעלויות השאילתות שמופעלות על התצוגה מחושבות על סמך העמודות שנבחרו בתצוגה, ולכן יכול להיות שגודל הבייטים שנסרקו יהיה גדול.

הזרמה לטבלאות מחולקות למחיצות לפי זמן

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

חלוקה למחיצות בזמן ההטמעה

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

נתונים חדשים שמגיעים ממוקמים באופן זמני במחיצה __UNPARTITIONED__ בזמן שהם במאגר הזמני של הנתונים. אם יש מספיק נתונים לא מחולקים, BigQuery מחלק את הנתונים למחיצה הנכונה. עם זאת, אין הסכם רמת שירות (SLA) לגבי משך הזמן שנדרש להעברת נתונים מהמחיצה __UNPARTITIONED__. אפשר להחריג נתונים מהמאגר הזמני של הנתונים שמוזרמים בשאילתה באמצעות סינון הערכים NULL מהמחיצה __UNPARTITIONED__ באמצעות אחד מהעמודות הווירטואליות (_PARTITIONTIME או _PARTITIONDATE, בהתאם לסוג הנתונים המועדף).

אם אתם מעבירים נתונים בסטרימינג לטבלה עם חלוקה למחיצות לפי יום, אתם יכולים לבטל את ההיסק של התאריך על ידי ציון מעצב מחיצות כחלק מ-Storage Write API ‏ (REST). כוללים את ה-decorator בפרמטר tableId. לדוגמה, אפשר להזרים למחיצה שמתאימה ל-2021-03-01 בטבלה table1 באמצעות קישוט המחיצה:

table1$20210301

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

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

לצורך בדיקה, אפשר להשתמש בכלי שורת הפקודה של BigQuery, בפקודת ה-CLI‏ bq insert. לדוגמה, הפקודה הבאה מעבירה שורה אחת למחיצה של התאריך 1 בינואר 2017 ($20170101) לטבלה עם מחיצות בשם mydataset.mytable:

echo '{"a":1, "b":2}' | bq insert 'mydataset.mytable$20170101'

חלוקה למחיצות (partitioning) לפי עמודה של יחידת זמן

אפשר להזרים נתונים לטבלה שמחולקת למחיצות לפי עמודה של DATE, DATETIME או TIMESTAMP, שערכיה הם בין 10 שנים בעבר לשנה אחת בעתיד. נתונים מחוץ לטווח הזה נדחים.

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

  • הערה: מחיצות יומיות עוברות עיבוד שונה ממחיצות שעתיות, חודשיות ושנתיות. רק נתונים מחוץ לטווח התאריכים (מ-7 הימים האחרונים עד 3 ימים קדימה) מחולצים למחיצה UNPARTITIONED, וממתינים לחלוקה מחדש למחיצות. לעומת זאת, בטבלה שמחולקת למחיצות לפי שעה, הנתונים תמיד מחולצים למחיצה UNPARTITIONED, ואחר כך מחולקים מחדש למחיצות.

יצירת טבלאות באופן אוטומטי באמצעות טבלאות תבנית

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

כדי להשתמש בטבלת תבנית דרך BigQuery API, מוסיפים פרמטר templateSuffix לבקשת Storage Write API (REST). בכלי שורת הפקודה של BigQuery, מוסיפים את הדגל template_suffix לפקודה insert. אם מערכת BigQuery מזהה פרמטר templateSuffix או דגל template_suffix, היא מתייחסת לטבלת היעד כתבנית בסיס. היא יוצרת טבלה חדשה עם אותה סכימה כמו הטבלה הממוקדת, ושם שכולל את הסיומת שצוינה:

 + 

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

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

שינוי סכימת הטבלה של התבנית

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

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

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

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

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

פרטי טבלת התבנית

ערך הסיומת של התבנית
הערך של templateSuffix (או --template_suffix) חייב להכיל רק אותיות (a-z, A-Z), מספרים (0-9) או קווים תחתונים (_). האורך המקסימלי של שם הטבלה והסיומת שלה ביחד הוא 1,024 תווים.
מכסה

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

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

משך החיים (TTL)

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

ביטול כפילויות

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

תצוגות

טבלת התבנית והטבלאות שנוצרו לא יכולות להיות תצוגות.

פתרון בעיות שקשורות להוספת שידורים חיים

בקטעים הבאים מוסבר איך לפתור בעיות שמתרחשות כשמעבירים נתונים ל-BigQuery באמצעות Storage Write API (REST). מידע נוסף על פתרון שגיאות שקשורות למכסות של הזנת זרם נתונים זמין במאמר בנושא שגיאות שקשורות למכסות של הזנת זרם נתונים.

קודי תגובת HTTP של כשל

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

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

קודי תגובת HTTP שמציינים הצלחה

גם אם מקבלים קוד תגובה של תגובת HTTP שמעיד על הצלחה, צריך לבדוק את המאפיין insertErrors של התגובה כדי לדעת אם הוספת השורות הצליחה, כי יכול להיות ש-BigQuery הצליח להוסיף רק חלק מהשורות. יכול להיות שתיתקלו באחד מהתרחישים הבאים:

  • כל השורות נוספו בהצלחה: אם המאפיין insertErrors הוא רשימה ריקה, כל השורות נוספו בהצלחה.
  • חלק מהשורות נוספו בהצלחה: למעט במקרים שבהם יש אי התאמה בסכימה באחת מהשורות, השורות שמצוינות במאפיין insertErrors לא נוספות, וכל שאר השורות נוספות בהצלחה. המאפיין errors מכיל מידע מפורט על הסיבה לכשל בכל שורה לא מוצלחת. המאפיין index מציין את אינדקס השורה (מבוסס-0) של הבקשה שהשגיאה רלוונטית לגביה.
  • לא הוכנסו שורות בהצלחה: אם ב-BigQuery מזוהה אי התאמה בסכימה בשורות נפרדות בבקשה, אף אחת מהשורות לא מוכנסת ומוחזרת רשומה של insertErrors לכל שורה, גם לשורות שלא הייתה בהן אי התאמה בסכימה. בשורה שלא הייתה בה אי התאמה לסכימה, המאפיין reason מוגדר לערך stopped, ואפשר לשלוח אותה מחדש כמו שהיא. בשורות שנכשלו מופיע מידע מפורט על אי ההתאמה לסכימה. מידע נוסף על סוגי מאגרי אחסון לפרוטוקולים ונתוני Arrow שנתמכים

שגיאות במטא-נתונים של הזנת זרם נתונים

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

דוגמאות לתרחישים:

  • שינויים בסכימה: שינוי הסכימה של טבלה שקיבלה לאחרונה הוספות של נתונים בזמן אמת עלול לגרום לתגובות עם שגיאות של חוסר התאמה בסכימה, כי יכול להיות שמערכת הסטרימינג לא תזהה את השינוי בסכימה באופן מיידי.
  • יצירה או מחיקה של טבלה: סטרימינג לטבלה שלא קיימת מחזיר וריאציה של תגובת notFound. יכול להיות שתוספות של נתונים בסטרימינג לא יזהו מיד טבלה שנוצרה בתגובה. באופן דומה, מחיקה או יצירה מחדש של טבלה יכולות ליצור תקופה שבה הוספות של נתונים בסטרימינג מועברות לטבלה הישנה. יכול להיות שהוספות הסטרימינג לא יופיעו בטבלה החדשה.
  • חיתוך טבלה: חיתוך של נתוני טבלה (באמצעות עבודת שאילתה שמשתמשת בערך writeDisposition של WRITE_TRUNCATE) עלול לגרום להשמטה של הוספות עוקבות במהלך תקופת העקביות.

נתונים חסרים או לא זמינים

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

שגיאות שקשורות למכסת הזנת זרם נתונים

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

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

הודעת שגיאה

אם השדה insertId ריק, יכול להיות שתופיע שגיאת המכסה הבאה:

מגבלת מכסה הודעת השגיאה
בייטים לשנייה לכל פרויקט הישות שלך עם gaia_id: GAIA_ID, פרויקט: PROJECT_ID באזור: REGION חרגה מהמכסה של בייטים להוספה לשנייה.

אם השדה insertId מאוכלס, יכולות להופיע שגיאות לגבי מכסת השימוש הבאות:

מגבלת מכסה הודעת השגיאה
שורות לשנייה לכל פרויקט הפרויקט שלך: PROJECT_ID ב-REGION חרג מהמכסה של הוספת שורות לשנייה בסטרימינג.
שורות לשנייה לכל טבלה הטבלה: TABLE_ID חרגה מהמכסה של הזנת זרם נתונים לשנייה.
בייטים לשנייה לכל טבלה הטבלה שלך: TABLE_ID חרגה מהמכסה של הזנת זרם נתונים בייטים לשנייה.

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

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

אבחון

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

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

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

SELECT
 start_timestamp,
 error_code,
 SUM(total_input_bytes) as sum_input_bytes,
 SUM(IF(error_code IN ('QUOTA_EXCEEDED', 'RATE_LIMIT_EXCEEDED'),
     total_requests, 0)) AS quota_error
FROM
 `region-REGION_NAME`.INFORMATION_SCHEMA.STREAMING_TIMELINE_BY_PROJECT
WHERE
  start_timestamp > TIMESTAMP_SUB(CURRENT_TIMESTAMP, INTERVAL 1 DAY)
GROUP BY
 start_timestamp,
 error_code
ORDER BY 1 DESC

רזולוציה

כדי לפתור את השגיאה שקשורה למכסת השימוש, צריך לבצע את הפעולות הבאות:

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

  • אם אתם לא משתמשים ב-insertId, או אם לא ניתן להסיר אותו, כדאי לעקוב אחרי התנועה של הסטרימינג במשך 24 שעות ולנתח את שגיאות המכסה:

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

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

    • אם אתם רואים שגיאות QUOTA_EXCEEDED או שתנועת הגולשים הכוללת חורגת באופן עקבי מ-80% מהמכסה, אתם יכולים לשלוח בקשה להגדלת המכסה. מידע נוסף זמין במאמר בנושא שליחת בקשה לשינוי המכסות.

    • אפשר גם להחליף את ההוספות של נתוני סטרימינג ב-Storage Write API החדש יותר, שכולל תפוקה גבוהה יותר, מחיר נמוך יותר ותכונות שימושיות רבות.