Custom connectors used to mean Scala or Java. Since 4.0 you subclass DataSource (name and schema) and DataSourceReader or DataSourceWriter (the work) in Python and register the class.
import json
from pyspark.sql.datasource import DataSource, DataSourceReader
class BookNestCatalog(DataSource):
@classmethod
def name(cls): return "booknest"
def schema(self): return "id INT, title STRING, price DOUBLE"
def reader(self, schema): return CatalogReader(self.options["path"])
class CatalogReader(DataSourceReader):
def __init__(self, path): self.path = path
def read(self, partition): # runs in a Python worker on an executor
for b in json.load(open(self.path))["books"]:
yield b["id"], b["title"], b["price"]
spark.dataSource.register(BookNestCatalog)
spark.read.format("booknest").load("data/raw/books.json").orderBy("price").show(3)Output
+---+--------------------+-----+ | id| title|price| +---+--------------------+-----+ | 1| The Quiet Harbor|14.99| | 5|The Clockmaker's ...| 16.2| | 4|Small Steps to Bi...|18.75| +---+--------------------+-----+ only showing top 3 rows
read() runs in Python workers once per partition (override partitions() to read in parallel), and rows return as Arrow 129 batches, so workers need pyarrow 129 : with the system python3 as worker this listing failed with No module named 'pyarrow' until PYSPARK_PYTHON pointed at the virtual environment.