התחברות למקורות נתונים
Managed Service for Apache Spark משתלב עם מסדי נתונים, מחסני נתונים ושירותי סטרימינג כדי לעבד נתונים ממערכות שונות. במאמר הזה נסביר איך לחבר אשכולות של Managed Service for Apache Spark למקורות נתונים בתוך Google Cloudומחוץ ל-Google Cloud.
לפני שמתחילים
- נכנסים לחשבון Google Cloud . אם אתם משתמשים חדשים ב- Google Cloud, צרו חשבון כדי שתוכלו להעריך את הביצועים של המוצרים שלנו בתרחישים מהעולם האמיתי. לקוחות חדשים מקבלים בחינם גם קרדיט בשווי 300$ להרצה, לבדיקה ולפריסה של עומסי העבודה.
-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that you have the permissions required to complete this guide.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Managed Service for Apache Spark, BigQuery, and Cloud Storage APIs.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.-
In the Google Cloud console, on the project selector page, select or create a Google Cloud project.
Roles required to select or create a project
- Select a project: Selecting a project doesn't require a specific IAM role—you can select any project that you've been granted a role on.
-
Create a project: To create a project, you need the Project Creator role
(
roles/resourcemanager.projectCreator), which contains theresourcemanager.projects.createpermission. Learn how to grant roles.
-
Verify that you have the permissions required to complete this guide.
-
Verify that billing is enabled for your Google Cloud project.
Enable the Managed Service for Apache Spark, BigQuery, and Cloud Storage APIs.
Roles required to enable APIs
To enable APIs, you need the
serviceusage.services.enablepermission. If you created the project, then you likely already have this permission through the Owner role (roles/owner). Otherwise, you can get this permission through the Service Usage Admin role (roles/serviceusage.serviceUsageAdmin). Learn how to grant roles.- יצירת אשכול של Managed Service for Apache Spark.
- יצירת קטגוריה של Cloud Storage
- חשוב לוודא שיש לכם את ההרשאות הנדרשות לגישה למקורות נתונים.
התפקידים הנדרשים
כדי להריץ את הדוגמאות בדף הזה, צריך תפקידים מסוימים ב-IAM. יכול להיות שהתפקידים האלה כבר הוקצו, בהתאם למדיניות הארגון. כדי לבדוק את התפקידים שהוקצו, ראו האם צריך להקצות תפקידים?.
כדי לקרוא הסבר על מתן תפקידים, ראו איך מנהלים את הגישה ברמת הפרויקט, התיקייה והארגון.
תפקידי משתמשים
כדי לקבל את ההרשאות שדרושות ליצירת אשכול של Managed Service for Apache Spark, צריך לבקש מהאדמין להקצות לכם את תפקידי ה-IAM הבאים:
- עריכת פרויקטים ב-Dataproc (
roles/dataproc.editor) בפרויקט - משתמש בחשבון שירות (
roles/iam.serviceAccountUser) בחשבון השירות של Compute Engine שמוגדר כברירת מחדל
תפקיד בחשבון שירות
כדי לוודא שלחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine יש את ההרשאות שנדרשות ליצירת אשכול של Managed Service for Apache Spark, צריך לבקש מהאדמין להקצות לחשבון השירות שמוגדר כברירת מחדל ב-Compute Engine את תפקיד ה-IAM Dataproc Worker (roles/dataproc.worker) בפרויקט.
חיבור לשירותים של Google Cloud
בקטע הזה נסביר איך להתחבר אל Managed Service for Apache Spark Google Cloud ומקורות נתונים חיצוניים.
חיבור ל-Cloud Storage
Managed Service for Apache Spark משתמש ב-Cloud Storage כמערכת הקבצים המבוזרת שלו כברירת מחדל. כדי לקרוא ולכתוב נתונים מקטגוריות של Cloud Storage, משתמשים בתחילית ה-URI gs://. אין צורך בהגדרות נוספות.
בדוגמה הבאה רואים איך לקרוא קובץ CSV מקטגוריה של Cloud Storage ואז לכתוב קובץ Parquet לקטגוריה.
# Read a CSV file from a Cloud Storage bucket into a Spark DataFrame
df = spark.read.format("csv") \
.option("header", "true") \
.option("inferSchema", "true") \
.load("gs://BUCKET_NAME/path/to/data.csv")
df.show()
# Write a Spark DataFrame to a Cloud Storage bucket
df.write.format("parquet") \
.mode("overwrite") \
.save("gs://BUCKET_NAME/path/to/output.parquet")
התחברות ל-BigQuery
מחבר BigQuery מותקן מראש ב-Managed Service for Apache Spark
2.1 ובגרסאות מאוחרות יותר של תמונות. המחבר מאפשר העברת נתונים בהיקף גדול בין Managed Service for Apache Spark לבין BigQuery.
בדוגמה הבאה מוצג איך לקרוא טבלת BigQuery לתוך Spark DataFrame, ואז להריץ שאילתה.
# Set up the BigQuery connector.
spark.conf.set("spark.sql.sources.provider", "bigquery")
spark.conf.set("viewsEnabled", "true")
spark.conf.set("materializationDataset", "YOUR_BQ_DATASET")
# Read a BigQuery table into a Spark DataFrame,
df = spark.read.format("bigquery") \
.option("table", "bigquery-public-data.samples.shakespeare") \
.load()
df.createOrReplaceTempView("shakespeare")
# Perform a query on the data,
word_counts = spark.sql("SELECT word, SUM(word_count) as count FROM shakespeare GROUP BY word ORDER BY count DESC")
word_counts.show()
חיבור למסדי נתונים חיצוניים באמצעות JDBC
אתם יכולים לחבר את Managed Service for Apache Spark למסד נתונים שמספק מנהל התקן Java Database Connectivity (JDBC), כמו PostgreSQL או MySQL, באמצעות מקור הנתונים המובנה של Spark JDBC.
הגדרת מנהל התקן JDBC כזמין
כדי להתחבר למסד נתונים באמצעות JDBC, צריך לוודא שקובץ ה-JAR של מנהל ההתקן זמין בנתיב המחלקה של Spark.
לכל משרה
מאחסנים את קובץ ה-JAR של הדרייבר בקטגוריה של Cloud Storage ומפנים אליו באמצעות הדגל --jars כששולחים עבודה. Spark מפיץ את קובץ ה-JAR לצמתים הדרושים למשימה. זו הגישה המומלצת לתלות ספציפית במשימה.
בכל האשכול
אם כל העבודות באשכול דורשות מנהל התקן ספציפי, אפשר להשתמש בפעולת אתחול כשיוצרים את האשכול כדי להעתיק את קובץ ה-JAR של מנהל ההתקן מ-Cloud Storage לספריית קובצי ה-JAR של Spark בכל צומת.
יוצרים סקריפט של פעולת אתחול.
#!/bin/bash set -e -x gsutil cp gs://YOUR_BUCKET_NAME/drivers/postgresql-42.7.3.jar /usr/lib/spark/jars/יוצרים את האשכול ומפנים לתסריט.
gcloud dataproc clusters create my-jdbc-cluster \ --region=REGION \ --image-version=2.2-debian11 \ --initialization-actions=gs://YOUR_BUCKET_NAME/init-actions/install-jdbc-driver.sh
קריאה ממסד נתונים של PostgreSQL
בדוגמה הזו נעשה שימוש בשיטה של כל משימה כדי להתחבר למסד נתונים של PostgreSQL.
מורידים את קובץ ה-JAR של מנהל ההתקן של PostgreSQL JDBC ומעלים אותו לקטגוריה של Cloud Storage.
שולחים את עבודת ה-PySpark ומפנים לנתיב של מנהל ההתקן ב-Cloud Storage באמצעות הדגל
--jars.gcloud dataproc jobs submit pyspark my_job.py \ --cluster=MY_CLUSTER \ --region=REGION \ --jars=gs://YOUR_BUCKET_NAME/drivers/postgresql-42.7.3.jarכדי לקרוא ולכתוב במסד הנתונים, משתמשים בקוד הבא בקובץ
job.py.from pyspark.sql import SparkSession def main(): # WARNING: Do not hardcode credentials. Use a service like Secret Manager # to handle sensitive credentials. db_properties = { "user": "USERNAME", "password": "PASSWORD", "driver": "org.postgresql.Driver" } jdbc_url = "jdbc:postgresql://DB_HOST:DB_PORT/DATABASE" spark = SparkSession.builder.appName("Postgres JDBC Example").getOrCreate() # Read data from a PostgreSQL table df = spark.read.jdbc(url=jdbc_url, table="public.my_table", properties=db_properties) df.printSchema() df.show() # Write data to a new PostgreSQL table (df.write .jdbc(url=jdbc_url, table="public.new_table", mode="overwrite", properties=db_properties)) if __name__ == "__main__": main()
חיבור למסדי נתונים חיצוניים באמצעות ODBC
כדי להתחבר למקורות נתונים שבהם מנהל התקן JDBC לא זמין, משתמשים במנהל התקן Open Database Connectivity (ODBC). משתמשים בפעולת אתחול כדי להתקין את מנהל ההתקן ואת התלות שלו בכל צומת כשיוצרים את האשכול.
בדוגמה הזו מוסבר איך להתחבר למופע של Microsoft SQL Server.
יוצרים את סקריפט פעולת האתחול.
- יוצרים סקריפט שמתקין את מנהל ה-ODBC של Microsoft SQL Server עבור Debian.
#!/bin/bash # Initialization action for installing MS SQL ODBC driver on the cluster. set -e -x # Install dependencies for the driver. apt-get update apt-get install -y --no-install-recommends curl gnupg unixodbc-dev # Add Microsoft's official repository. curl https://packages.microsoft.com/keys/microsoft.asc | apt-key add - curl https://packages.microsoft.com/config/debian/11/prod.list > /etc/apt/sources.list.d/mssql-release.list # Install the driver. apt-get update ACCEPT_EULA=Y apt-get install -y msodbcsql18 # Install the pyodbc library for Python. pip install pyodbc ``` 1. Upload the script to a Cloud Storage bucket. ```shell gsutil cp install-odbc-mssql.sh gs://YOUR_BUCKET_NAME/init-actions/ ```יוצרים אשכול של Managed Service for Apache Spark שמריץ את פעולת האתחול.
כשיוצרים את האשכול, מצביעים על סקריפט פעולת האתחול ב-Cloud Storage.
gcloud dataproc clusters create my-odbc-cluster \ --region=REGION \ --image-version=2.2-debian11 \ --initialization-actions=gs://BUCKET_NAME/init-actions/install-odbc-mssql.shחיבור וביצוע שאילתות באמצעות PySpark
אחרי שהאשכול פועל, מנהל התקן ה-ODBC והספרייה
pyodbcזמינים. הקוד הבא משתמש ב-pyodbcבצומת דרייבר כדי לאחזר נתונים ל-Pandas DataFrame, ואז ממיר אותם ל-Spark DataFrame לעיבוד מבוזר.import pyodbc import pandas as pd from pyspark.sql import SparkSession def get_sql_data(connection_string, query): """ Connects to the database using pyodbc, executes a query, and returns the result as a Pandas DataFrame. """ cnxn = pyodbc.connect(connection_string) pdf = pd.read_sql(query, cnxn) cnxn.close() return pdf def main(): spark = SparkSession.builder.appName("ODBC Example").getOrCreate() # WARNING: Do not hardcode credentials. Use Secret Manager. server = 'SERVER.database.windows.net' database = 'DATABASE' username = 'USERNAME' password = 'PASSWORD' driver = '{ODBC Driver 18 for SQL Server}' # The driver name must match the installed driver. # Create the connection string. connection_string = f'DRIVER={driver};SERVER=tcp:{server},1433;DATABASE={database};UID={username};PWD={password};Encrypt=yes;TrustServerCertificate=no;Connection Timeout=30;' sql_query = "SELECT * FROM Sales.Customer" # Fetch data as a Pandas DataFrame. pandas_df = get_sql_data(connection_string, sql_query) # Convert the Pandas DataFrame to a Spark DataFrame for distributed processing. spark_df = spark.createDataFrame(pandas_df) print("Successfully read data from SQL Server:") spark_df.printSchema() spark_df.show(5) # You can now perform distributed operations on the Spark DataFrame. print(f"Total number of customers: {spark_df.count()}") if __name__ == "__main__": main()
שיטות מומלצות
בקטע הזה מתוארות שיטות מומלצות לניהול פרטי כניסה ותלות.
ניהול פרטי כניסה
אל תקודדו סיסמאות או סודות אחרים בקוד. אפשר להשתמש ב-Secret Manager כדי לאחסן פרטי כניסה בצורה מאובטחת ולגשת אליהם ממשימות של Managed Service for Apache Spark.
ניהול יחסי תלות
אפשר להשתמש בפעולות אתחול כדי להתקין מנהלי התקנים וספריות כשיוצרים אשכול. בסביבות מורכבות, כדאי ליצור תמונה מותאמת אישית של Managed Service for Apache Spark עם כל יחסי התלות שכבר מותקנים. שימוש באימג' בהתאמה אישית מקצר את זמן ההפעלה של האשכול ומבטיח עקביות.
לפני שמשתמשים בסקריפטים של פעולות אתחול בסביבת ייצור, צריך להעתיק אותם ממקורות ציבוריים לקטגוריה משלכם ב-Cloud Storage עם ניהול גרסאות. אל תפנו ישירות לסקריפטים ציבוריים, כי הם יכולים להשתנות בלי הודעה מראש.