bigquery-streaming-connector
This is a library for bigquery streaming connector for pyspark structured streaming
The underlying connector uses bigquery storage api services to pull bigquery table data at scale using spark workers.
The storage api services is cheaper and faster than the traditional Bigquery query api services enabling faster & cheaper Bigquery migration incrementally in a continuous fashion.
Pre-requisite
Need spark 4.0.0 or Databricks runtime version 15.3 & above.
pip install bigquery-spark-streaming-connector
Pyspark usage:
from streaming_connector import bq_stream_register
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 max 1000
.load()
## The above will ingest table data incrementally using the provided timestamp based field and latest value is checkpointed using offset semantics.
## Without the incremental input field full table ingestion is done.
## The service_account_json files need to be available to every spark executor workers in the '/home' folder using --files /home/<file_name>.json options or using init script (using Databricks spark)
(query.writeStream.trigger(processingTime='30 seconds')
.option("checkpointLocation", "checkpoint_path")
.foreachBatch(writeToTable) #your target table write function
.start()
)
Release files for bigquery-spark-streaming-connector 0.7.0
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| bigquery_spark_streaming_connector-0.7.0-py3-none-any.whl | Python 3 | none | any | Details |
Release files / bigquery_spark_streaming_connector-0.7.0-py3-none-any.whl
| Download URL | bigquery_spark_streaming_connector-0.7.0-py3-none-any.whl |
|---|---|
| Size | 6.4 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
f6399c06852e6fbd0085ca1c6e0de599b36fff259351f8038616e0f3421e83db
|
|
BLAKE2b-256 checksum How to use checksums |
ac59adcc6a407dbd9c98ec26a47c624572527eb38a6ea8fbc3eced9f954cd8d0
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
No |
| Uploaded via |
twine/5.1.1 CPython/3.12.2
|