beam-nuggets
Collection of transforms for the Apache beam python SDK.
Decision gist · record as of 2026-08-14
Yes, if you need Beam-to-database or Beam-to-Kafka connectors and can tolerate dormant maintenance. The package has low install friction and no known vulnerabilities, but verify compatibility with your Beam and driver versions first. If you require active support or need to integrate with very recent Beam or database driver releases, consider forking or maintaining a local patch.AI-flagged interpretation of the facts on this page — verify before relying
Before you install
- Requires Apache Beam to be installed; database drivers (pg8000 for PostgreSQL, PyMySQL for MySQL) must be available for the target database.
- Low friction installation with a pure-Python wheel.
- Maintenance is dormant—last release was 2021-09-12 and no commits since 2023-12-07—so expect no active bug fixes or updates for newer Beam or database driver versions.
License · maintenance · safety
permissive license (permissive) — MIT license (permissive) allows commercial use, modification, and distribution with minimal restrictions.
last release 2021-09-12 (1797 days) · last repo commit 2023-12-07 · 90 stars
0 known vulnerabilities (OSV.dev, 2026-08-14) · 76,789 downloads/mo, #14,590 on PyPI
Alternatives
Verify before relying
pip install beam-nuggets
import apache_beam as beam
from beam_nuggets.io import relational_db
source_config = relational_db.SourceConfiguration(
drivername='sqlite',
database='/tmp/test.sqlite'
)
table_config = relational_db.TableConfiguration(name='data')
with beam.Pipeline() as p:
p | beam.Create([{'id': 1}]) | relational_db.Write(
source_config=source_config,
table_config=table_config
)- Compatibility with recent Apache Beam versions (last tested against unknown version as of 2021).
- Whether Kafka transforms work with modern kafka-python API and broker versions.
- Support for Python versions beyond 3.x (exact minor versions unspecified in metadata).
What it is and what it does
beam-nuggets is a collection of custom transforms that extend Apache Beam's Python SDK with connectors for relational databases, Kafka, and CSV files. It wraps SQLAlchemy to provide ReadFromDB and Write transforms that work with any SQLAlchemy-supported database (PostgreSQL, MySQL, SQLite tested), plus KafkaConsume and KafkaProduce for Kafka integration and CSV reading. It also includes utility transforms for parsing JSON, selecting from nested dictionaries, and assigning unique IDs.
The package is most useful for data pipelines that need to ingest from or load into SQL databases or Kafka topics within a Beam pipeline. It abstracts away the boilerplate of configuring database connections and serialization, letting you focus on pipeline logic. However, the project is dormant—last updated in September 2021—so it may not work with recent versions of Beam, database drivers, or Kafka clients without manual fixes.
Use it for
- Read records from a PostgreSQL or MySQL table and process them in a Beam pipeline.
- Write transformed data from a Beam pipeline directly into a relational database table.
- Consume messages from a Kafka topic and process them in a Beam pipeline.
- Produce processed records to a Kafka topic from a Beam pipeline.
- Parse CSV files or JSON data as part of a Beam data processing workflow.
Worth the install?
AI-flagged interpretation of the facts on this page. Verify before relying on it.
Yes, if you need Beam-to-database or Beam-to-Kafka connectors and can tolerate dormant maintenance.
The package has low install friction and no known vulnerabilities, but verify compatibility with your Beam and driver versions first. If you require active support or need to integrate with very recent Beam or database driver releases, consider forking or maintaining a local patch.
Install
beam-nuggets on PyPI
Before you install
Low friction installation with a pure-Python wheel. Maintenance is dormant—last release was 2021-09-12 and no commits since 2023-12-07—so expect no active bug fixes or updates for newer Beam or database driver versions.
Requires Apache Beam to be installed; database drivers (pg8000 for PostgreSQL, PyMySQL for MySQL) must be available for the target database.
License in practice
MIT license (permissive) allows commercial use, modification, and distribution with minimal restrictions.
Quickstart
pip install beam-nuggets
import apache_beam as beam
from beam_nuggets.io import relational_db
source_config = relational_db.SourceConfiguration(
drivername='sqlite',
database='/tmp/test.sqlite'
)
table_config = relational_db.TableConfiguration(name='data')
with beam.Pipeline() as p:
p | beam.Create([{'id': 1}]) | relational_db.Write(
source_config=source_config,
table_config=table_config
)
Verify before relying
- Compatibility with recent Apache Beam versions (last tested against unknown version as of 2021).
- Whether Kafka transforms work with modern kafka-python API and broker versions.
- Support for Python versions beyond 3.x (exact minor versions unspecified in metadata).
Package facts
| License | permissive license permissive |
| Python support | Not specified |
| Install friction | Low. Pure-Python wheel |
| Runtime dependencies | 6 packagesapache-beamSQLAlchemysqlalchemy-utilspg8000PyMySQLkafka-python |
| Maintenance | Dormant 1,797 days since the last release |
| Last repo commit | |
| First released | |
| Downloads | 76,789 / month, #14,590 on PyPI 30-day window, as of 2026-08-14 |
| Known vulnerabilities | None known OSV.dev, checked 2026-08-14 |
| Classifiers | License :: OSI Approved :: MIT LicenseOperating System :: OS IndependentProgramming Language :: PythonProgramming Language :: Python :: 3 |
Evidence: beam_nuggets-0.18.1-py3-none-any.whl
Tags
Let your AI agent find packages like this
Example. Real query, live index.
You found this page by searching. An agent finds it by wishing: SkillFed indexes 14,416 PyPI packages by what they can do, searchable in plain language.
wish › “apache beam database transforms”
- beam-nuggetsProvides Apache Beam transforms for reading and writing relational…
- apache-airflow-providers-apache-beamIntegrates Apache Beam data processing pipelines into Apache Airflow…
- apache-beamApache Beam is a unified framework for defining and executing batch…
Give your agent the search over MCP, or paste the wish link into any chat.
More Distributed Computing packages
gRPC Python is an HTTP/2-based RPC framework that enables you to define and call remote procedures across network boundaries using protocol buffers for serialization.
Install it if you need RPC communication in a distributed system or are integrating with existing gRPC services.
execnet lets you spawn and communicate with Python interpreters across local processes, remote hosts, and different platforms, using a simple API for task distribution and inter-process messaging.
However, the aging maintenance status (275 days since last release) means you should verify it meets your concurrency and performance needs before committing to a…
Cloudpickle extends Python's standard pickle module to serialize lambda functions, interactively-defined functions and classes, and other constructs that the default pickle cannot handle, making it suitable for cluster computing and remote code execution.
Install it if you need to serialize lambda functions, interactively-defined code, or non-standard Python constructs for cluster computing or distributed execution.
Provides a unified, open()-compatible Python API for streaming large files from remote storage (S3, GCS, Azure, HDFS, SFTP, HTTP) and local filesystems, with transparent compression support.
Install it if you work with large files on cloud storage or remote systems and want to avoid writing boilerplate around multiple SDKs.
Portalocker provides cross-platform file locking with support for exclusive and shared locks, plus Redis-based distributed locks and process-aware PID file locking.
Install it if you need file or process coordination; the optional extras (pywin32, redis) are only required for specific lock types.
Ray is a distributed computing framework that scales Python applications from a single machine to multi-node clusters, providing abstractions for parallel tasks, stateful actors, and shared objects.
See also apache-beam · apache-airflow-providers-apache-beam · llama-index-storage-kvstore-postgres · agate-sql · pangres · quixstreams · kafka · databases · dataset · tortoise-orm