Kommentar
Åtkomst till den här sidan kräver auktorisering. Du kan prova att logga in eller ändra kataloger.
Åtkomst till den här sidan kräver auktorisering. Du kan prova att ändra kataloger.
PySpark-anpassade datakällor skapas med hjälp av Python-api:et (PySpark) DataSource, som gör det möjligt att läsa från anpassade datakällor och skriva till anpassade datamottagare i Apache Spark med hjälp av Python. Du kan använda anpassade PySpark-datakällor för att definiera anpassade anslutningar till datasystem och implementera ytterligare funktioner för att skapa återanvändbara datakällor.
Spark inkluderar inbyggt stöd för standardformat som Delta, Iceberg, Parquet, JSON, CSV och JDBC, men inte för många andra system som REST API:er, Google Sheets, Hugging Face-dataset eller proprietära interna tjänster. Python DataSource API fyller detta gap: du bygger kopplingar till dessa system i ren Python, utan JVM-baserad kopplingsutveckling, och använder dem som vilken inbyggd Spark-datakälla som helst, inklusive i Spark SQL.
Kommentar
Anpassade PySpark-datakällor kräver Databricks Runtime 15.4 LTS och senare, eller serverlös miljöversion 2.
DataSource-klass
PySpark DataSource är en basklass som tillhandahåller metoder för att skapa dataläsare och författare.
Implementera datakällans underklass
Beroende på ditt användningsfall måste följande implementeras av alla underklasser för att göra en datakälla antingen läsbar, skrivbar eller både och:
| Egenskap eller metod | beskrivning |
|---|---|
name |
Obligatoriskt. Namnet på datakällan |
schema |
Obligatoriskt. Schemat för datakällan som ska läsas eller skrivas |
reader() |
Måste returnera en DataSourceReader för att göra datakällan läsbar i batchläge. |
writer() |
Måste returnera en DataSourceWriter för att göra datasänkan skrivbar (satsvis) |
streamReader() eller simpleStreamReader() |
Måste returnera en DataSourceStreamReader för att göra dataströmmen läsbar (direktuppspelning) |
streamWriter() |
Måste returnera en DataSourceStreamWriter för att göra dataströmmen skrivbar (strömningsbar) |
Kommentar
De användardefinierade DataSourcemetoderna , DataSourceReader, DataSourceWriter, DataSourceStreamReader, DataSourceStreamWriteroch deras metoder måste vara serialiserbara. Med andra ord måste de vara en ordlista eller kapslad ordlista som innehåller en primitiv typ.
Registrera datakällan
När du har implementerat gränssnittet måste du registrera det, sedan kan du läsa in eller på annat sätt använda det som visas i följande exempel:
# Register the data source
spark.dataSource.register(MyDataSourceClass)
# Read from a custom data source
spark.read.format("my_datasource_name").load().show()
Exempel 1: Skapa en PySpark DataSource för batchfråga
Om du vill demonstrera PySpark DataSource-läsarfunktioner skapar du en datakälla som genererar exempeldata med hjälp av Python-paketet faker . Mer information om fakerfinns i Faker-dokumentationen.
faker Installera paketet med följande kommando:
%pip install faker
Steg 1: Implementera läsaren för en batchfråga
Implementera först läsarlogik för att generera exempeldata. Använd det installerade faker biblioteket för att fylla i varje fält i schemat.
class FakeDataSourceReader(DataSourceReader):
def __init__(self, schema, options):
self.schema: StructType = schema
self.options = options
def read(self, partition):
# Library imports must be within the method.
from faker import Faker
fake = Faker()
# Every value in this `self.options` dictionary is a string.
num_rows = int(self.options.get("numRows", 3))
for _ in range(num_rows):
row = []
for field in self.schema.fields:
value = getattr(fake, field.name)()
row.append(value)
yield tuple(row)
Steg 2: Definiera exemplet DataSource
Definiera sedan din nya PySpark DataSource som en underklass av DataSource med namn, schema och läsare. Metoden reader() måste definieras för att läsa från en datakälla i en batchfråga.
from pyspark.sql.datasource import DataSource, DataSourceReader
from pyspark.sql.types import StructType
class FakeDataSource(DataSource):
"""
An example data source for batch query using the `faker` library.
"""
@classmethod
def name(cls):
return "fake"
def schema(self):
return "name string, date string, zipcode string, state string"
def reader(self, schema: StructType):
return FakeDataSourceReader(schema, self.options)
Steg 3: Registrera och använda exempeldatakällan
Om du vill använda datakällan registrerar du den. Som standard har de FakeDataSource tre raderna och schemat innehåller följande string fält: name, date, zipcode, . state Följande exempel registrerar, läser in och matar ut exempeldatakällan med standardvärdena:
spark.dataSource.register(FakeDataSource)
spark.read.format("fake").load().show()
+-----------------+----------+-------+----------+
| name| date|zipcode| state|
+-----------------+----------+-------+----------+
|Christine Sampson|1979-04-24| 79766| Colorado|
| Shelby Cox|2011-08-05| 24596| Florida|
| Amanda Robinson|2019-01-06| 57395|Washington|
+-----------------+----------+-------+----------+
Endast string fält stöds, men du kan ange ett schema med alla fält som motsvarar faker paketleverantörernas fält för att generera slumpmässiga data för testning och utveckling. I följande exempel läses datakällan in med fälten name och company:
spark.read.format("fake").schema("name string, company string").load().show()
+---------------------+--------------+
|name |company |
+---------------------+--------------+
|Tanner Brennan |Adams Group |
|Leslie Maxwell |Santiago Group|
|Mrs. Jacqueline Brown|Maynard Inc |
+---------------------+--------------+
Om du vill läsa in datakällan med ett anpassat antal rader anger du alternativet numRows . I följande exempel anges 5 rader:
spark.read.format("fake").option("numRows", 5).load().show()
+--------------+----------+-------+------------+
| name| date|zipcode| state|
+--------------+----------+-------+------------+
| Pam Mitchell|1988-10-20| 23788| Tennessee|
|Melissa Turner|1996-06-14| 30851| Nevada|
| Brian Ramsey|2021-08-21| 55277| Washington|
| Caitlin Reed|1983-06-22| 89813|Pennsylvania|
| Douglas James|2007-01-18| 46226| Alabama|
+--------------+----------+-------+------------+
Exempel 2: Skriv till en anpassad datasänkare i en batchfråga
För att demonstrera skrivfunktionerna i PySpark DataSource, skapa en datakälla som skriver varje partition i en DataFrame till en fil och sedan skriver en sammanfattande markeringsfil när jobbet slutförs.
Steg 1: Implementera skrivmodulen för en batchfråga
Först, implementera skrivarlogiken. Varje exekutor anropar write() en gång per partition. När alla skrivuppgifter har lyckats anropar drivrutinen commit(). Om någon uppgift misslyckas anropar drivrutinen abort() i stället.
from dataclasses import dataclass
from pyspark.sql.datasource import DataSourceWriter, WriterCommitMessage
@dataclass
class SimpleCommitMessage(WriterCommitMessage):
partition_id: int
count: int
class FakeDataSourceWriter(DataSourceWriter):
def __init__(self, options):
self.path = options.get("path")
assert self.path is not None
def write(self, iterator):
"""
Writes the rows in a partition to a file, then returns a commit message with the row count. Library imports must be within the method.
"""
import json
import os
from pyspark import TaskContext
# Runs on an executor, so create the output directory on the local node.
os.makedirs(self.path, exist_ok=True)
partition_id = TaskContext.get().partitionId()
count = 0
with open(os.path.join(self.path, f"part-{partition_id}.json"), "w") as file:
for row in iterator:
file.write(json.dumps(row.asDict()) + "\n")
count += 1
return SimpleCommitMessage(partition_id=partition_id, count=count)
def commit(self, messages):
"""
Runs on the driver after all write tasks succeed. Writes a summary of the write to a marker file.
"""
import json
import os
# Runs on the driver, so create the output directory on the local node.
os.makedirs(self.path, exist_ok=True)
total_rows = sum(message.count for message in messages if message is not None)
with open(os.path.join(self.path, "_SUCCESS"), "w") as file:
file.write(json.dumps({"partitions": len(messages), "rows": total_rows}))
def abort(self, messages):
"""
Runs on the driver if any write task fails. Use it to clean up partial output.
"""
import os
# Runs on the driver, so create the output directory on the local node.
os.makedirs(self.path, exist_ok=True)
with open(os.path.join(self.path, "_FAILED"), "w") as file:
file.write("write job aborted")
Steg 2: Definiera en skrivbar DataSource
Definiera sedan en DataSource underklass som implementerar writer(). Argumentet overwrite är True när skrivläget är overwrite och False när det är append.
from pyspark.sql.datasource import DataSource
from pyspark.sql.types import StructType
class FakeSinkDataSource(DataSource):
"""
An example writable data source that saves rows to files.
"""
@classmethod
def name(cls):
return "fakesink"
def schema(self):
return "name string, date string, zipcode string, state string"
def writer(self, schema: StructType, overwrite: bool):
return FakeDataSourceWriter(self.options)
Steg 3: Registrera och skriv till datasamlaren
Om du vill använda datakällan registrerar du den. Skriv sedan en DataFrame till den genom att skicka det korta namnet till format() och en utdatakatalog till alternativet path . Detta exempel skriver till en sökväg i en Unity Catalog-volym. Ersätt <catalog>, <schema>, och <volume> med en befintlig volym.
spark.dataSource.register(FakeSinkDataSource)
output_path = "/Volumes/<catalog>/<schema>/<volume>/fakesink"
df = spark.range(3).selectExpr(
"cast(id as string) as name",
"'2025-01-01' as date",
"'12345' as zipcode",
"'California' as state",
)
df.write.format("fakesink").mode("append").option("path", output_path).save()
Kommentar
Antalet utdatafiler är lika med antalet partitioner i DataFrame, inte antalet rader. Antalet partitioner kommer från klustrets standardparallellism, så med fler partitioner än rader får vissa partitioner inga rader och producerar tomma filer.
Exempel 3: Skapa en PySpark GitHub DataSource med varianter
För att demonstrera användningen av varianter i en PySpark DataSource skapar det här exemplet en datakälla som läser pull-begäranden från GitHub.
Kommentar
Varianter stöds med anpassade PySpark-datakällor i Databricks Runtime 17.1 och senare.
Information om varianter finns i Frågevariantdata.
Steg 1: Implementera läsaren för att hämta pull-begäranden
Implementera först läsarlogik för att hämta pull-begäranden från den angivna GitHub-lagringsplatsen.
class GithubVariantPullRequestReader(DataSourceReader):
def __init__(self, options):
self.token = options.get("token")
self.repo = options.get("path")
if self.repo is None:
raise Exception(f"Must specify a repo in `.load()` method.")
# Every value in this `self.options` dictionary is a string.
self.num_rows = int(options.get("numRows", 10))
def read(self, partition):
header = {
"Accept": "application/vnd.github+json",
}
if self.token is not None:
header["Authorization"] = f"Bearer {self.token}"
url = f"https://api.github.com/repos/{self.repo}/pulls"
response = requests.get(url, headers=header)
response.raise_for_status()
prs = response.json()
for pr in prs[:self.num_rows]:
yield Row(
id = pr.get("number"),
title = pr.get("title"),
user = VariantVal.parseJson(json.dumps(pr.get("user"))),
created_at = pr.get("created_at"),
updated_at = pr.get("updated_at")
)
Steg 2: Definiera GitHub DataSource
Definiera sedan din nya PySpark GitHub DataSource som en underklass av DataSource med ett namn, schema och en metod reader(). Schemat innehåller följande fält: id, title, user, , created_at. updated_at Fältet user definieras som en variant.
import json
import requests
from pyspark.sql import Row
from pyspark.sql.datasource import DataSource, DataSourceReader
from pyspark.sql.types import VariantVal
class GithubVariantDataSource(DataSource):
@classmethod
def name(self):
return "githubVariant"
def schema(self):
return "id int, title string, user variant, created_at string, updated_at string"
def reader(self, schema):
return GithubVariantPullRequestReader(self.options)
Steg 3: Registrera och använda datakällan
Om du vill använda datakällan registrerar du den. Följande exempel registreras och läser sedan in datakällan och matar ut tre rader av GitHub-lagringsplatsens PR-data:
spark.dataSource.register(GithubVariantDataSource)
spark.read.format("githubVariant").option("numRows", 3).load("apache/spark").display()
+---------+-----------------------------------------------------+---------------------+----------------------+----------------------+
| id | title | user | created_at | updated_at |
+---------+---------------------------------------------------- +---------------------+----------------------+----------------------+
| 51293 |[SPARK-52586][SQL] Introduce AnyTimeType | {"avatar_url":...} | 2025-06-26T09:20:59Z | 2025-06-26T15:22:39Z |
| 51292 |[WIP][PYTHON] Arrow UDF for aggregation | {"avatar_url":...} | 2025-06-26T07:52:27Z | 2025-06-26T07:52:37Z |
| 51290 |[SPARK-50686][SQL] Hash to sort aggregation fallback | {"avatar_url":...} | 2025-06-26T06:19:58Z | 2025-06-26T06:20:07Z |
+---------+-----------------------------------------------------+---------------------+----------------------+----------------------+
Exempel 4: Skapa PySpark DataSource för strömmande läsning och skrivning
Skapa en exempeldatakälla som genererar två rader i varje mikrobatch med hjälp av faker Python-paketet för att demonstrera strömläsare och skrivarfunktioner i PySpark DataSource. Mer information om fakerfinns i Faker-dokumentationen.
faker Installera paketet med följande kommando:
%pip install faker
Steg 1: Implementera strömläsaren
Implementera först exempelläsaren för strömmande data som genererar två rader i varje mikrobatch. Du kan implementera DataSourceStreamReader, eller om datakällan har lågt dataflöde och inte kräver partitionering, kan du implementera SimpleDataSourceStreamReader i stället. Antingen simpleStreamReader() eller streamReader() måste implementeras och simpleStreamReader() anropas endast när streamReader() inte har implementerats.
Implementering av DataSourceStreamReader
Instansen streamReader har en heltalsförskjutning som ökar med 2 i varje mikrobatch, implementerad med gränssnittet DataSourceStreamReader.
from pyspark.sql.datasource import InputPartition
from typing import Iterator, Tuple
import os
import json
class RangePartition(InputPartition):
def __init__(self, start, end):
self.start = start
self.end = end
class FakeStreamReader(DataSourceStreamReader):
def __init__(self, schema, options):
self.current = 0
def initialOffset(self) -> dict:
"""
Returns the initial start offset of the reader.
"""
return {"offset": 0}
def latestOffset(self) -> dict:
"""
Returns the current latest offset that the next microbatch will read to.
"""
self.current += 2
return {"offset": self.current}
def partitions(self, start: dict, end: dict):
"""
Plans the partitioning of the current microbatch defined by start and end offset. It
needs to return a sequence of :class:`InputPartition` objects.
"""
return [RangePartition(start["offset"], end["offset"])]
def commit(self, end: dict):
"""
This is invoked when the query has finished processing data before end offset. This
can be used to clean up the resource.
"""
pass
def read(self, partition) -> Iterator[Tuple]:
"""
Takes a partition as an input and reads an iterator of tuples from the data source.
"""
start, end = partition.start, partition.end
for i in range(start, end):
yield (i, str(i))
Implementering av SimpleDataSourceStreamReader
Instansen SimpleStreamReader är samma som den FakeStreamReader instans som genererar två rader i varje batch, men implementeras med SimpleDataSourceStreamReader gränssnittet utan partitionering.
class SimpleStreamReader(SimpleDataSourceStreamReader):
def initialOffset(self):
"""
Returns the initial start offset of the reader.
"""
return {"offset": 0}
def read(self, start: dict) -> (Iterator[Tuple], dict):
"""
Takes start offset as an input, then returns an iterator of tuples and the start offset of the next read.
"""
start_idx = start["offset"]
it = iter([(i,) for i in range(start_idx, start_idx + 2)])
return (it, {"offset": start_idx + 2})
def readBetweenOffsets(self, start: dict, end: dict) -> Iterator[Tuple]:
"""
Takes start and end offset as inputs, then reads an iterator of data deterministically.
This is called when the query replays batches during restart or after a failure.
"""
start_idx = start["offset"]
end_idx = end["offset"]
return iter([(i,) for i in range(start_idx, end_idx)])
def commit(self, end):
"""
This is invoked when the query has finished processing data before end offset. This can be used to clean up resources.
"""
pass
Steg 2: Implementera strömskrivaren
Implementera sedan strömningsskrivaren. Den här dataskrivaren för strömning skriver metadatainformationen för varje mikrobatch till en lokal filväg.
from pyspark.sql.datasource import DataSourceStreamWriter, WriterCommitMessage
class SimpleCommitMessage(WriterCommitMessage):
def __init__(self, partition_id: int, count: int):
self.partition_id = partition_id
self.count = count
class FakeStreamWriter(DataSourceStreamWriter):
def __init__(self, options):
self.options = options
self.path = self.options.get("path")
assert self.path is not None
def write(self, iterator):
"""
Writes the data and then returns the commit message for that partition. Library imports must be within the method.
"""
from pyspark import TaskContext
context = TaskContext.get()
partition_id = context.partitionId()
cnt = 0
for row in iterator:
cnt += 1
return SimpleCommitMessage(partition_id=partition_id, count=cnt)
def commit(self, messages, batchId) -> None:
"""
Receives a sequence of :class:`WriterCommitMessage` when all write tasks have succeeded, then decides what to do with it.
In this FakeStreamWriter, the metadata of the microbatch(number of rows and partitions) is written into a JSON file inside commit().
"""
status = dict(num_partitions=len(messages), rows=sum(m.count for m in messages))
with open(os.path.join(self.path, f"{batchId}.json"), "a") as file:
file.write(json.dumps(status) + "\n")
def abort(self, messages, batchId) -> None:
"""
Receives a sequence of :class:`WriterCommitMessage` from successful tasks when some other tasks have failed, then decides what to do with it.
In this FakeStreamWriter, a failure message is written into a text file inside abort().
"""
with open(os.path.join(self.path, f"{batchId}.txt"), "w") as file:
file.write(f"failed in batch {batchId}")
Steg 3: Definiera exemplet DataSource
Definiera nu din nya PySpark DataSource som en underklass med DataSource ett namn, schema och metoder streamReader() och streamWriter().
from pyspark.sql.datasource import DataSource, DataSourceStreamReader, SimpleDataSourceStreamReader, DataSourceStreamWriter
from pyspark.sql.types import StructType
class FakeStreamDataSource(DataSource):
"""
An example data source for streaming read and write using the `faker` library.
"""
@classmethod
def name(cls):
return "fakestream"
def schema(self):
return "name string, state string"
def streamReader(self, schema: StructType):
return FakeStreamReader(schema, self.options)
# If you don't need partitioning, you can implement the simpleStreamReader method instead of streamReader.
# def simpleStreamReader(self, schema: StructType):
# return SimpleStreamReader()
def streamWriter(self, schema: StructType, overwrite: bool):
return FakeStreamWriter(self.options)
Steg 4: Registrera och använda exempeldatakällan
Om du vill använda datakällan registrerar du den. När den har registrerats kan du använda den i strömmande frågor som källa eller mottagare genom att skicka ett kort namn eller fullständigt namn till format(). I följande exempel registreras datakällan, och sedan startas en fråga som läser från exempeldatakällan och skickar utdata till konsolen:
spark.dataSource.register(FakeStreamDataSource)
query = spark.readStream.format("fakestream").load().writeStream.format("console").start()
Alternativt använder följande kod exempelströmmen som en mottagare och anger en utdatasökväg:
spark.dataSource.register(FakeStreamDataSource)
# Make sure the output directory exists and is writable
output_path = "/output_path"
dbutils.fs.mkdirs(output_path)
checkpoint_path = "/output_path/checkpoint"
query = (
spark.readStream
.format("fakestream")
.load()
.writeStream
.format("fakestream")
.option("path", output_path)
.option("checkpointLocation", checkpoint_path)
.start()
)
Exempel 5: Skapa en Google BigQuery-strömningskoppling
I följande exempel visas hur du skapar en anpassad anslutningsapp för direktuppspelning för Google BigQuery (BQ) med hjälp av en PySpark DataSource. Databricks tillhandahåller en Spark-anslutning för BigQuery-batchinmatning, och Lakehouse Federation kan också fjärransluta till alla BigQuery-datauppsättningar och hämta data genom att skapa en extern katalog, men inget av dem stöder fullt ut arbetsflöden för inkrementell eller kontinuerlig strömning. Den här anslutningsappen möjliggör stegvis inkrementell datamigrering och nästan realtidsmigrering från BigQuery-tabeller som matas av strömmande källor med beständiga kontrollpunkter.
Den här anpassade anslutningsappen har följande funktioner:
- Kompatibel med Structured Streaming och Lakeflow-pipelines.
- Stöder inkrementell postspårning och kontinuerlig strömningsinmatning och följer semantik för strukturerad direktuppspelning.
- Använder BigQuery Storage API med ett RPC-baserat protokoll för snabbare och billigare dataöverföring.
- Skriver migrerade tabeller direkt till Unity Catalog.
- Hanterar kontrollpunkter automatiskt med hjälp av ett datum- eller tidsstämpelbaserat inkrementellt fält.
- Stöder batchinmatning med
Trigger.AvailableNow(). - Kräver ingen mellanliggande molnlagring.
- Serialiserar BigQuery-dataöverföring med pil- eller Avro-format.
- Hanterar automatisk parallellitet och distribuerar arbete mellan Spark-arbetare baserat på datavolym.
- Passar för raw- och bronslagermigrering från BigQuery, med stöd för silver- och guldlagermigreringar med scd-mönster av typ 1 eller typ 2.
Förutsättningar
Innan du implementerar den anpassade anslutningsappen installerar du de paket som krävs:
%pip install faker google.cloud google.cloud.bigquery google.cloud.bigquery_storage
Steg 1: Implementera strömläsaren
Implementera först strömmande dataläsare. Underklassen DataSourceStreamReader måste implementera följande metoder:
initialOffset(self) -> dictlatestOffset(self) -> dictpartitions(self, start: dict, end: dict) -> Sequence[InputPartition]read(self, partition: InputPartition) -> Union[Iterator[Tuple], Iterator[Row]]commit(self, end: dict) -> Nonestop(self) -> None
Mer information om varje metod finns i Metoder.
import os
from pyspark.sql.datasource import DataSourceStreamReader, InputPartition
from pyspark.sql.datasource import DataSourceStreamWriter
from pyspark.sql import Row
from pyspark.sql import SparkSession
from pyspark.sql.datasource import DataSource
from pathlib import Path
from pyarrow.lib import TimestampScalar
from datetime import datetime
from typing import Iterator, Tuple, Any, Dict, List, Sequence
from google.cloud.bigquery_storage import BigQueryReadClient, ReadSession
from google.cloud import bigquery_storage
import pandas
import datetime
import uuid
import time, logging
start_time = time.time()
class RangePartition(InputPartition):
def __init__(self, session: ReadSession, stream_idx: int):
self.session = session
self.stream_idx = stream_idx
class BQStreamReader(DataSourceStreamReader):
def __init__(self, schema, options):
self.project_id = options.get("project_id")
self.dataset = options.get("dataset")
self.table = options.get("table")
self.json_auth_file = "/home/"+options.get("service_auth_json_file_name")
self.max_parallel_conn = options.get("max_parallel_conn", 1000)
self.incremental_checkpoint_field = options.get("incremental_checkpoint_field", "")
self.last_offset = None
def initialOffset(self) -> dict:
"""
Returns the initial start offset of the reader.
"""
from datetime import datetime
logging.info("Inside initialOffset!!!!!")
# self.increment_latest_vals.append(datetime.strptime('1900-01-01 23:57:12', "%Y-%m-%d %H:%M:%S"))
self.last_offset = '1900-01-01 23:57:12'
return {"offset": str(self.last_offset)}
def latestOffset(self):
"""
Returns the current latest offset that the next microbatch will read to.
"""
from datetime import datetime
from google.cloud import bigquery
if (self.last_offset is None):
self.last_offset = '1900-01-01 23:57:12'
client = bigquery.Client.from_service_account_json(self.json_auth_file)
# max_offset=start["offset"]
logging.info(f"************************last_offset: {self.last_offset}***********************")
f_sql_str = ''
for x_str in self.incremental_checkpoint_field.strip().split(","):
f_sql_str += f"{x_str}>'{self.last_offset}' or "
f_sql_str = f_sql_str[:-3]
job_query = client.query(
f"select max({self.incremental_checkpoint_field}) from {self.project_id}.{self.dataset}.{self.table} where {f_sql_str}")
for query in job_query.result():
max_res = query[0]
if (str(max_res).lower() != 'none'):
return {"offset": str(max_res)}
return {"offset": str(self.last_offset)}
def partitions(self, start: dict, end: dict) -> Sequence[InputPartition]:
"""
Plans the partitioning of the current microbatch defined by start and end offset. It
needs to return a sequence of :class:`InputPartition` objects.
"""
if (self.last_offset is None):
self.last_offset = end['offset']
os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = self.json_auth_file
# project_id = self.auth_project_id
client = BigQueryReadClient()
# This example reads baby name data from the public datasets.
table = "projects/{}/datasets/{}/tables/{}".format(
self.project_id, self.dataset, self.table
)
requested_session = bigquery_storage.ReadSession()
requested_session.table = table
if (self.incremental_checkpoint_field != ''):
start_offset = start["offset"]
end_offset = end["offset"]
f_sql_str = ''
for x_str in self.incremental_checkpoint_field.strip().split(","):
f_sql_str += f"({x_str}>'{start_offset}' and {x_str}<='{end_offset}') or "
f_sql_str = f_sql_str[:-3]
requested_session.read_options.row_restriction = f"{f_sql_str}"
# This example leverages Apache Avro.
requested_session.data_format = bigquery_storage.DataFormat.AVRO
parent = "projects/{}".format(self.project_id)
session = client.create_read_session(
request={
"parent": parent,
"read_session": requested_session,
"max_stream_count": int(self.max_parallel_conn),
},
)
self.last_offset = end['offset']
return [RangePartition(session, i) for i in range(len(session.streams))]
def read(self, partition) -> Iterator[List]:
"""
Takes a partition as an input and reads an iterator of tuples from the data source.
"""
from datetime import datetime
session = partition.session
stream_idx = partition.stream_idx
os.environ['GOOGLE_APPLICATION_CREDENTIALS'] = self.json_auth_file
client_1 = BigQueryReadClient()
# requested_session.read_options.selected_fields = ["census_tract", "clearance_date", "clearance_status"]
reader = client_1.read_rows(session.streams[stream_idx].name)
reader_iter = []
for message in reader.rows():
reader_iter_in = []
for k, v in message.items():
reader_iter_in.append(v)
# yield(reader_iter)
reader_iter.append(reader_iter_in)
# yield (message['hash'], message['size'], message['virtual_size'], message['version'])
# self.increment_latest_vals.append(max_incr_val)
return iter(reader_iter)
def commit(self, end):
"""
This is invoked when the query has finished processing data before end offset. This
can be used to clean up the resource.
"""
pass
Steg 2: Definiera DataSource
Definiera sedan den anpassade datakällan. Underklassen DataSource måste implementera följande metoder:
name(cls) -> strschema(self) -> Union[StructType, str]
Mer information om varje metod finns i Metoder.
from pyspark.sql.datasource import DataSource
from pyspark.sql.types import StructType
from google.cloud import bigquery
class BQStreamDataSource(DataSource):
"""
An example data source for streaming data from a public API containing users' comments.
"""
@classmethod
def name(cls):
return "bigquery-streaming"
def schema(self):
type_map = {'integer': 'long', 'float': 'double', 'record': 'string'}
json_auth_file = "/home/" + self.options.get("service_auth_json_file_name")
client = bigquery.Client.from_service_account_json(json_auth_file)
table_ref = self.options.get("project_id") + '.' + self.options.get("dataset") + '.' + self.options.get("table")
table = client.get_table(table_ref)
original_schema = table.schema
result = []
for schema in original_schema:
col_attr_name = schema.name
if (schema.mode != 'REPEATED'):
col_attr_type = type_map.get(schema.field_type.lower(), schema.field_type.lower())
else:
col_attr_type = f"array<{type_map.get(schema.field_type.lower(), schema.field_type.lower())}>"
result.append(col_attr_name + " " + col_attr_type)
return ",".join(result)
# return "census_tract double,clearance_date string,clearance_status string"
def streamReader(self, schema: StructType):
return BQStreamReader(schema, self.options)
Steg 3: Konfigurera och starta strömningsfrågan
Registrera slutligen kontakten, konfigurera och starta sedan strömningsfrågan:
spark.dataSource.register(BQStreamDataSource)
# Ingests table data incrementally using the provided timestamp-based field.
# The latest value is checkpointed using offset semantics.
# Without the incremental input field, full table ingestion is performed.
# Service account JSON files must be available to every Spark executor worker
# in the /home folder using --files /home/<file_name>.json or an init script.
query = (
spark.readStream.format("bigquery-streaming")
.option("project_id", <bq_project_id>)
.option("incremental_checkpoint_field", <table_incremental_ts_based_col>)
.option("dataset", <bq_dataset_name>)
.option("table", <bq_table_name>)
.option("service_auth_json_file_name", <service_account_json_file_name>)
.option("max_parallel_conn", <max_parallel_threads_to_pull_data>) # defaults to max 1000
.load()
)
(
query.writeStream.trigger(processingTime="30 seconds")
.option("checkpointLocation", "checkpoint_path")
.foreachBatch(writeToTable) # your target table write function
.start()
)
Körningsordning
Funktionens körningsordning för den anpassade dataströmmen beskrivs nedan.
För inläsning av Spark Stream DataFrame:
name(cls)
schema()
För mikrobatch (n) vid start av en ny förfrågan eller när du startar om en befintlig förfrågan (ny eller befintlig kontrollpunkt):
partitions(end_offset, end_offset) # loads the last saved offset from the checkpoint at query restart
latestOffset()
partitions(start_offset, end_offset) # plans partitions and distributes to Python workers
read() # user’s source read definition, runs on each Python worker
commit()
För nästa (n+1) mikrobatch av en fråga som körs på en befintlig kontrollpunkt:
latestOffset()
partitions(start_offset, end_offset)
read()
commit()
Kommentar
Funktionen latestOffset samordnar kontrollpunkter. Dela en kontrollpunktsvariabel av en primitiv typ mellan funktioner och returnera den som en ordlista. Till exempel: return {"offset": str(self.last_offset)}
Exempel 6: Autentisera med ett externt API
Det här exemplet visar hur du autentiserar en PySpark-datakälla med ett externt HTTP-API med hjälp av en HTTP-anslutning i Unity Catalog, så att datakällkoden aldrig innehåller hårdkodade token eller autentiseringsuppgifter.
Kommentar
Inmatning av HTTP-anslutningsautentiseringsuppgifter i Unity Catalog kräver Databricks Runtime 18.1 eller senare.
Steg 1: Skapa en HTTP-anslutning
Innan du implementerar datakällan skapar du en HTTP-anslutning med namnet my_weather_api i Unity Catalog och ger användare eller grupper MANAGE behörighet till den. Endast användare med behörighet för MANAGE anslutningen kan utlösa inmatning av autentiseringsuppgifter.
Lagra API-token som en Databricks-hemlighet och referera till den secret med funktionen i stället för att ange literaltoken, så att autentiseringsuppgifterna aldrig visas i anslutningsdefinitionen.
CREATE CONNECTION my_weather_api TYPE HTTP
OPTIONS (
host 'https://api.openweathermap.org',
base_path '/data/2.5',
bearer_token secret('my_secret_scope', 'weather_api_token')
);
GRANT MANAGE ON CONNECTION my_weather_api TO `user@example.com`;
Steg 2: Implementera läsaren för en batchfråga
Implementera sedan läsarlogik för att hämta data från REST-API:et. Läsaren läser de inmatade hostvärdena , base_pathoch bearer_token från dess alternativ, så inga autentiseringsuppgifter visas i koden.
from pyspark.sql.datasource import DataSource, DataSourceReader, InputPartition
from urllib.parse import quote
import urllib.error
import urllib.request
import json
class WeatherApiReader(DataSourceReader):
def __init__(self, options):
self.host = options["host"]
self.base_path = options["base_path"]
self.token = options["bearer_token"]
# Every value in this `options` dictionary is a string.
self.cities = options.get("cities", "Seattle,Portland,Denver").split(",")
def partitions(self):
return [InputPartition(city.strip()) for city in self.cities]
def read(self, partition):
city = partition.value
# URL-encode the city so names with spaces or non-ASCII characters (for example, "New York" or "São Paulo") produce a valid query string.
url = f"{self.host}{self.base_path}/weather?q={quote(city)}&units=metric"
req = urllib.request.Request(url)
req.add_header("Authorization", f"Bearer {self.token}")
try:
# Set a timeout so a slow or unresponsive API surfaces a controlled error instead of hanging the Spark task.
with urllib.request.urlopen(req, timeout=30) as resp:
data = json.loads(resp.read().decode())
except (urllib.error.URLError, TimeoutError) as e:
raise RuntimeError(f"Weather API request failed for {city}: {e}")
# Validate the response shape before indexing so an error payload raises a clear message instead of a KeyError.
try:
main = data["main"]
weather = data["weather"][0]
except (KeyError, IndexError, TypeError):
raise RuntimeError(f"Unexpected weather API response for {city}: {data}")
yield (city, main["temp"], main["humidity"], weather["description"])
Steg 3: Definiera exemplet DataSource
Definiera nu din nya PySpark DataSource som en underklass av DataSource med namn, schema och läsare.
class WeatherApiSource(DataSource):
def __init__(self, options):
self.options = options
@classmethod
def name(cls):
return "weather_api"
def schema(self):
return "city STRING, temperature DOUBLE, humidity INT, description STRING"
def reader(self, schema):
return WeatherApiReader(self.options)
Steg 4: Registrera och använda datakällan
Om du vill använda datakällan registrerar du den. Referera sedan till HTTP-anslutningen för Unity Catalog med alternativet databricks.connection . Spark-drivrutinen hämtar automatiskt kortlivade OAuth2-autentiseringsuppgifter från Unity Catalog och matar in dem (till exempel bearer_token, hostoch base_path) i datakällsalternativkartan. Autentiseringsnycklar som matas in av Unity Catalog kan inte åsidosättas och alternativ som är globalt blockerade, till exempel host och port, förblir blockerade och kan inte anges av användare.
spark.dataSource.register(WeatherApiSource)
df = (
spark.read.format("weather_api")
.option("databricks.connection", "my_weather_api") # Unity Catalog injects host, base_path, bearer_token
.option("cities", "Seattle,Portland,Denver") # user-defined option passes through
.load()
)
df.show()
Det här exemplet implementerar endast batchläsningar. Samma databricks.connection alternativ gäller även för läsningar och skrivningar för direktuppspelning när datakällan implementerar motsvarande metoder (streamReader eller simpleStreamReader för direktuppspelningsläsningar och writer eller streamWriter för skrivningar).
Ytterligare resurser
Apache Spark-communityn underhåller exempelkontakter som du kan använda som referenser när du bygger dina egna datakällor. Dessa arkiv underhålls av communityn och stöds inte av Databricks:
- pyspark-data-sources: En samling exempel på PySpark anpassade datakällakopplingar.
- pyspark_huggingface: En anpassad datakällanslutning för att läsa datauppsättningar från Hugging Face.
Felsökning
Om du får följande felmeddelande betyder det att din dator inte stöder anpassade PySpark-datakällor. Du måste använda Databricks Runtime 15.2 eller senare.
Error: [UNSUPPORTED_FEATURE.PYTHON_DATA_SOURCE] The feature is not supported: Python data sources. SQLSTATE: 0A000