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):

Reply via email to