For AI agents: the complete documentation index is available at https://docs.dataplatform.ovh.net/ja/llms.txt, the full documentation bundle is available at https://docs.dataplatform.ovh.net/ja/llms-full.txt, and this page is available as Markdown at https://docs.dataplatform.ovh.net/ja/developers-python-sdk-quick-start-spark.md.
  • 🇯🇵 日本語
  • PySparkのユースケース

    Data Platform Python SDKを使用した5つのPySparkユースケース:バケットの読み取りからオブジェクトストレージへの書き込みまで

    目的

    このガイドでは、Data Platform Python SDKを使用した5つのPySparkユースケースについて説明します。バケットからの読み取り、SQLを使用したデータベースのクエリ、DataFrameの抽出、バケットまたはオブジェクトストレージへの書き込みです。

    ユースケース1:バケットから読み取り、データセットに書き込み

    このサンプルでは、Data Platform Bucketbucket_testからデータを取得し、Lakehouse Managerdefault_datasetデータセットに挿入する方法を示します。

    接続文字列の詳細については、Data Platform Buckets ConnectorLakehouse Manager Dataset Connectorを参照してください。

    from logging import getLogger
    from forepaas.dwh import connect, update_metas
    from pyspark import SparkContext
    from pyspark.sql import SQLContext
    
    logger = getLogger(__name__)
    
    cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
    cn_default = connect("dwh/default_dataset/")
    
    # Spark compatible connector extract_dataframe function returns Spark DataFrame
    spark_df = cn_source.extract_dataframe()
    logger.notice(f"CSV Columns: {list(spark_df.columns)}")
    
    # insert_dataframe uses Spark DataFrame as well
    cn_default.insert_dataframe("chicago_calendar_full_copy", spark_df)
    
    # At the end you can run update_metas() it will update the metas for all tables, so from lakehouse manager, you will see the correct number of rows
    update_metas()
    

    ユースケース2:データベースをSQLでクエリ

    このサンプルでは、Data Platform Connector objectget_spark_options()およびget_spark_context()メソッドを使用してデータベースをクエリする方法を紹介します。

    from logging import getLogger
    from forepaas.dwh import connect
    from pyspark import SparkContext
    from pyspark.sql import SQLContext
    
    logger = getLogger(__name__)
    
    # get_spark_context() and get_spark_options() are available for database connectors (snowflake, postgresql, mysql)
    cn_default = connect("dwh/default_dataset/")
    sc_default = cn_default.get_spark_context()
    so_default = cn_default.get_spark_options()
    
    # Depends on Snowflake or PostgreSQL
    sql_driver = "net.snowflake.spark.snowflake" # "jdbc" or "net.snowflake.spark.snowflake"
    
    spark_default = SQLContext(sc_default)
    sql = "select * from chicago_calendar_full"
    spark_df = spark_default.read.format(sql_driver).options(**so_default).option("query", sql).load()
    
    logger.notice(f"SQL Columns: {list(spark_df.columns)}")
    Info

    このユースケースはMySQL、PostgreSQL、Snowflakeのみで動作します

    ユースケース3:PySpark DataFrameを抽出(オプション付き)

    このサンプルでは、extract_dataframe()メソッドを使用してSpark DataFrameを素早く抽出する方法と、get_spark_url()get_spark_context()、およびget_spark_session()メソッドを使用して手動で行う方法を示します。

    from logging import getLogger
    from forepaas.dwh import connect
    from pyspark import SparkContext
    from pyspark.sql import SQLContext
    
    logger = getLogger(__name__)
    
    # Get all tables from default_dataset
    cn_default = connect("dwh/default_dataset/")
    
    # If needed you can print all tables
    # logger.info(cn_default.list())
    
    cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
    
    # Manual override extract_dataframe file options
    spark_df = cn_source.extract_dataframe(options)
    logger.notice(f"CSV1 Columns: {list(spark_df.columns)}")
    
    # Manual read from file
    # get_spark_url(), get_spark_context() and get_spark_session() are available for s3 / buckets connectors
    spark_session = cn_source.get_spark_session()
    
    # getting stations_rides.csv under bucket buc_test
    url = cn_source.get_spark_url("", "stations_rides.csv", bucket="buc_test")
    logger.notice(f"SparkURL: {url}")
    
    # Use format(file_suffix) for other files, check spark documentation for more information
    options= {"encoding": "utf-8", "sep": ";", "header": True}
    spark_df = spark_session.read.format("csv").options(**options).load(url)
    
    logger.notice(f"CSV2 Columns: {list(spark_df.columns)}")

    ユースケース4:別のバケットに書き込み

    このサンプルでは、get_spark_url()メソッドを使用して、1つのData Platform Bucketから別のバケットに書き込む方法を示します。

    from logging import getLogger
    from forepaas.dwh import connect
    from pyspark import SparkContext
    from pyspark.sql import SQLContext
    
    logger = getLogger(__name__)
    
    cn_source = connect("dwh/bucket_test/stations_rides.csv")
    spark_df = cn_source.extract_dataframe()
    
    url_dst = cn_source.get_spark_url("", "stations_rides_copy.csv", bucket="test2")
    logger.notice(f"SparkURL Dest: {url_dst}")
    spark_df.write.format("csv").options(**options).save(url_dst)
    
    url_dst = cn_source.get_spark_url("", "stations_rides_copy.parquet", bucket="test3")
    logger.notice(f"SparkURL Dest: {url_dst}")
    spark_df.write.format("parquet").save(url_dst)

    ユースケース5:insert_dataframeを使用してオブジェクトストレージに書き込み

    insert_dataframe()関数は、前のユースケースの操作を簡素化します。 これは、互換性のあるオブジェクトストレージタイプのPySparkコネクタ(現在のData Platform Buckets、S31、およびAzure Blob Storage)で利用可能です。

    from logging import getLogger
    from forepaas.dwh import connect
    
    logger = getLogger(__name__)
    
    cn_source = connect("dwh/bucket_test/chicago_calendar_full.csv")
    spark_df = cn_source.extract_dataframe()
    
    # cn_dest: Data Platform Buckets, S3, Azure Blob Storage 
    cn_dest = connect("dwh/dest_bucket/")
    
    # Insert using custom type
    params={"type": "csv"}
    # Insert using default type inferred from file name and default platform options
    # Raises an exception if file suffix not in ["csv", "json", "parquet"]
    # Destination path will be destination where you will find a .csv file under the configured path of the source
    cn_dest.insert_dataframe("destination", spark_df)
    
    # Insert using custom type
    params = {"type": "parquet"}
    # Reads from params.type, if not provided and no suffix in file name, an exception will be raised
    cn_dest.insert_dataframe("destination", spark_df, params)
    
    # Insert to specified absolute path
    params = {"path": "output/destination", "type": "csv"}
    # Destination file will be output/destination.csv
    cn_dest.insert_dataframe("", spark_df, params)
    
    # Insert with custom options
    params = {"write_options": {"sep": ",", "header": False}, "type": "csv"}
    cn_dest.insert_dataframe("destination", spark_df, params)
    
    # Current default platform options:
    # CSV: {"encoding": "utf-8", "sep": ";", "header": True}
    Info

    PySparkでPySpark DataFrameをファイルシステム(例:S3)に保存すると、PySparkは分散処理を行うため、単一のファイルの代わりにフォルダーを作成します。このフォルダー内には、DataFrameのパーティションを表す複数のpart-*.csvファイルが生成され、ファイルの数はDataFrameのパーティション数に依存します。また、成功した書き込み操作を示す空の_SUCCESSファイルも作成されます。

    さらに詳しく

    当社のソリューションを実装するためのトレーニングや技術サポートが必要な場合は、営業担当者にお問い合わせください、またはこのリンクをクリックして見積もりを依頼し、当社のプロフェッショナルサービスの専門家にプロジェクトのカスタム分析を依頼してください。

    Data Platformを構築するチームと直接やり取りし、質問をする、フィードバックを提供する、専用のDiscordチャネルに参加してください。

    OVHcloudサービスについてサポートが必要な場合は、ヘルプセンターでリクエストを作成してください。

    ユーザーコミュニティに参加してください。

    1: S3はAmazon Technologies, Inc.の商標です。OVHcloudのサービスはAmazon Technologies, Inc.によってスポンサー、承認、または提携されていません。