Skip to content

Commit 309f5fd

Browse files
committed
Refactor code structure for improved readability and maintainability
1 parent 0dfdf46 commit 309f5fd

4 files changed

Lines changed: 515 additions & 518 deletions

File tree

sdks/python/apache_beam/runners/dataflow/internal/apiclient.py

Lines changed: 8 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -28,23 +28,28 @@
2828

2929
import ast
3030
import codecs
31+
from functools import partial
3132
import getpass
3233
import hashlib
3334
import io
3435
import json
3536
import logging
3637
import os
3738
import random
38-
import re
3939
import string
40+
41+
from packaging import version
42+
import re
4043
import sys
4144
import time
4245
import traceback
4346
import warnings
4447
from copy import copy
4548
from datetime import datetime
4649
from datetime import timezone
47-
from functools import partial
50+
51+
from apitools.base.py import encoding
52+
from apitools.base.py import exceptions
4853

4954
from apache_beam import version as beam_version
5055
from apache_beam.internal.gcp.auth import get_service_credentials
@@ -68,11 +73,8 @@
6873
from apache_beam.transforms import cy_combiners
6974
from apache_beam.transforms.display import DisplayData
7075
from apache_beam.transforms.environments import is_apache_beam_container
71-
from apache_beam.utils import proto_utils
7276
from apache_beam.utils import retry
73-
from apitools.base.py import encoding
74-
from apitools.base.py import exceptions
75-
from packaging import version
77+
from apache_beam.utils import proto_utils
7678

7779
# Environment version information. It is passed to the service during a
7880
# a job submission and is used by the service to establish what features
@@ -203,11 +205,6 @@ def __init__(
203205
pool.diskSizeGb = self.worker_options.disk_size_gb
204206
if self.worker_options.disk_type:
205207
pool.diskType = self.worker_options.disk_type
206-
if self.worker_options.disk_provisioned_iops is not None:
207-
pool.diskProvisionedIops = self.worker_options.disk_provisioned_iops
208-
if self.worker_options.disk_provisioned_throughput_mibps is not None:
209-
pool.diskProvisionedThroughputMibps = (
210-
self.worker_options.disk_provisioned_throughput_mibps)
211208
if self.worker_options.zone:
212209
pool.zone = self.worker_options.zone
213210
if self.worker_options.network:

sdks/python/apache_beam/runners/dataflow/internal/apiclient_test.py

Lines changed: 0 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -128,25 +128,6 @@ def test_set_subnetwork(self):
128128
env.proto.workerPools[0].subnetwork,
129129
'/regions/MY/subnetworks/SUBNETWORK')
130130

131-
def test_set_disk_provisioning_options(self):
132-
pipeline_options = PipelineOptions([
133-
'--disk_provisioned_iops',
134-
'4000',
135-
'--disk_provisioned_throughput_mibps',
136-
'200',
137-
'--temp_location',
138-
'gs://any-location/temp',
139-
])
140-
env = apiclient.Environment(
141-
[], # packages
142-
pipeline_options,
143-
'2.0.0', # any environment version
144-
FAKE_PIPELINE_URL,
145-
)
146-
self.assertEqual(env.proto.workerPools[0].diskProvisionedIops, 4000)
147-
self.assertEqual(
148-
env.proto.workerPools[0].diskProvisionedThroughputMibps, 200)
149-
150131
def test_flexrs_blank(self):
151132
pipeline_options = PipelineOptions(
152133
['--temp_location', 'gs://any-location/temp'])

0 commit comments

Comments
 (0)