This is an automated email from the ASF dual-hosted git repository.
shunping pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git
The following commit(s) were added to refs/heads/master by this push:
new db27d09e5ec [Python] Make read and write buffer size configurable in
gcsio (#40184)
db27d09e5ec is described below
commit db27d09e5ec8aaff12f1ee57fad4c031ae3b9b5a
Author: Shunping Huang <[email protected]>
AuthorDate: Mon Sep 21 17:49:37 2026 -0400
[Python] Make read and write buffer size configurable in gcsio (#40184)
---
sdks/python/apache_beam/io/gcp/gcsio.py | 33 ++++++++++++++++++--
.../python/apache_beam/options/pipeline_options.py | 35 ++++++++++++++++++++++
2 files changed, 65 insertions(+), 3 deletions(-)
diff --git a/sdks/python/apache_beam/io/gcp/gcsio.py
b/sdks/python/apache_beam/io/gcp/gcsio.py
index e7372e8231e..a673845811d 100644
--- a/sdks/python/apache_beam/io/gcp/gcsio.py
+++ b/sdks/python/apache_beam/io/gcp/gcsio.py
@@ -67,6 +67,11 @@ __all__ = ['GcsIO', 'create_storage_client']
_LOGGER = logging.getLogger(__name__)
DEFAULT_READ_BUFFER_SIZE = 16 * 1024 * 1024
+DEFAULT_WRITE_BUFFER_SIZE = 16 * 1024 * 1024
+
+# Writes are performed as resumable uploads, which require the chunk size to be
+# a multiple of 256 KiB.
+WRITE_BUFFER_SIZE_MULTIPLE = 256 * 1024
# Maximum number of operations permitted in GcsIO.copy_batch() and
# GcsIO.delete_batch().
@@ -240,6 +245,14 @@ class GcsIO(object):
gcsio_retry.get_retry(pipeline_options) if GCS_INSTALLED else None)
self._use_blob_generation = getattr(
google_cloud_options, 'enable_gcsio_blob_generation', False)
+ self._read_buffer_size = getattr(
+ google_cloud_options, 'gcs_read_buffer_size_bytes', None)
+ if self._read_buffer_size is None:
+ self._read_buffer_size = DEFAULT_READ_BUFFER_SIZE
+ self._write_buffer_size = getattr(
+ google_cloud_options, 'gcs_write_buffer_size_bytes', None)
+ if self._write_buffer_size is None:
+ self._write_buffer_size = DEFAULT_WRITE_BUFFER_SIZE
def get_project_number(self, bucket):
if bucket not in self.bucket_to_project_number:
@@ -286,15 +299,23 @@ class GcsIO(object):
self,
filename,
mode='r',
- read_buffer_size=DEFAULT_READ_BUFFER_SIZE,
- mime_type='application/octet-stream'):
+ read_buffer_size=None,
+ mime_type='application/octet-stream',
+ write_buffer_size=None):
"""Open a GCS file path for reading or writing.
Args:
filename (str): GCS file path in the form ``gs://<bucket>/<object>``.
mode (str): ``'r'`` for reading or ``'w'`` for writing.
read_buffer_size (int): Buffer size to use during read operations.
+ Defaults to the value of the ``--gcs_read_buffer_size_bytes``
+ pipeline option, or ``DEFAULT_READ_BUFFER_SIZE`` when that option is
+ not set.
mime_type (str): Mime type to set for write operations.
+ write_buffer_size (int): Buffer size to use during write operations.
+ Must be a multiple of 256 KiB. Defaults to the value of the
+ ``--gcs_write_buffer_size_bytes`` pipeline option, or
+ ``DEFAULT_WRITE_BUFFER_SIZE`` when that option is not set.
Returns:
GCS file object.
@@ -302,6 +323,11 @@ class GcsIO(object):
Raises:
ValueError: Invalid open file mode.
"""
+ if read_buffer_size is None:
+ read_buffer_size = self._read_buffer_size
+ if write_buffer_size is None:
+ write_buffer_size = self._write_buffer_size
+
bucket_name, blob_name = parse_gcs_path(filename)
bucket = self.client.bucket(bucket_name)
@@ -317,6 +343,7 @@ class GcsIO(object):
return BeamBlobWriter(
blob,
mime_type,
+ chunk_size=write_buffer_size,
enable_write_bucket_metric=self.enable_write_bucket_metric,
retry=self._storage_client_retry)
else:
@@ -749,7 +776,7 @@ class BeamBlobWriter(BlobWriter):
self,
blob,
content_type,
- chunk_size=16 * 1024 * 1024,
+ chunk_size=DEFAULT_WRITE_BUFFER_SIZE,
ignore_flush=True,
enable_write_bucket_metric=False,
retry=DEFAULT_RETRY):
diff --git a/sdks/python/apache_beam/options/pipeline_options.py
b/sdks/python/apache_beam/options/pipeline_options.py
index 98e3dea38d2..6443094fda1 100644
--- a/sdks/python/apache_beam/options/pipeline_options.py
+++ b/sdks/python/apache_beam/options/pipeline_options.py
@@ -1205,6 +1205,23 @@ class GoogleCloudOptions(PipelineOptions):
'Entries are key value pairs separated by = '
'(e.g. --gcs_custom_audit_entry key=value) or a JSON string '
'(e.g. --gcs_custom_audit_entries=\'{ "user": "test", "id": "12"
}\').')
+ parser.add_argument(
+ '--gcs_read_buffer_size_bytes',
+ type=int,
+ default=None,
+ help='Size in bytes of the buffer used when reading from GCS. A '
+ 'larger buffer reduces the number of requests sent to GCS at the '
+ 'cost of more memory per reader. When unset, the GCS client in Beam '
+ 'uses its default buffer size (16 MiB).')
+ parser.add_argument(
+ '--gcs_write_buffer_size_bytes',
+ type=int,
+ default=None,
+ help='Size in bytes of the buffer used when writing to GCS. Must be '
+ 'a multiple of 256 KiB, since writes are performed as resumable '
+ 'uploads. A larger buffer reduces the number of requests sent to GCS '
+ 'at the cost of more memory per writer. When unset, the GCS client '
+ 'in Beam uses its default buffer size (16 MiB).')
def _create_default_gcs_bucket(self):
try:
@@ -1314,6 +1331,24 @@ class GoogleCloudOptions(PipelineOptions):
validator.validate_repeatable_argument_passed_as_list(
self, 'dataflow_service_options'))
+ if (self.gcs_read_buffer_size_bytes is not None and
+ self.gcs_read_buffer_size_bytes <= 0):
+ errors.append(
+ '--gcs_read_buffer_size_bytes must be a positive number of bytes, '
+ 'got %s.' % self.gcs_read_buffer_size_bytes)
+
+ if self.gcs_write_buffer_size_bytes is not None:
+ # GCS resumable uploads require the chunk size to be a multiple of
+ # 256 KiB. Checking here avoids a failure deep inside the GCS client
+ # on the first flush.
+ write_buffer_size_multiple = 256 * 1024
+ if (self.gcs_write_buffer_size_bytes <= 0 or
+ self.gcs_write_buffer_size_bytes % write_buffer_size_multiple != 0):
+ errors.append(
+ '--gcs_write_buffer_size_bytes must be a positive multiple of '
+ '%d bytes, got %s.' %
+ (write_buffer_size_multiple, self.gcs_write_buffer_size_bytes))
+
return errors
def get_cloud_profiler_service_name(self):