Skip to content

Commit 41194d4

Browse files
committed
remove Go sdk changes
1 parent f6a490a commit 41194d4

3 files changed

Lines changed: 7 additions & 20 deletions

File tree

sdks/go/pkg/beam/runners/dataflow/dataflow.go

Lines changed: 0 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -63,8 +63,6 @@ var (
6363
maxNumWorkers = flag.Int64("max_num_workers", 0, "Maximum number of workers during scaling (optional).")
6464
diskSizeGb = flag.Int64("disk_size_gb", 0, "Size of root disk for VMs, in GB (optional).")
6565
diskType = flag.String("disk_type", "", "Type of root disk for VMs (optional).")
66-
diskProvisionedIOPS = flag.Int64("disk_provisioned_iops", 0, "Provisioned IOPS for the root disk for VMs (optional).")
67-
diskProvisionedThroughputMibps = flag.Int64("disk_provisioned_throughput_mibps", 0, "Provisioned throughput for the root disk for VMs (optional).")
6866
autoscalingAlgorithm = flag.String("autoscaling_algorithm", "", "Autoscaling mode to use (optional).")
6967
zone = flag.String("zone", "", "GCP zone (optional)")
7068
kmsKey = flag.String("dataflow_kms_key", "", "The Cloud KMS key identifier used to encrypt data at rest (optional).")
@@ -117,8 +115,6 @@ var flagFilter = map[string]bool{
117115
"max_num_workers": true,
118116
"disk_size_gb": true,
119117
"disk_type": true,
120-
"disk_provisioned_iops": true,
121-
"disk_provisioned_throughput_mibps": true,
122118
"autoscaling_algorithm": true,
123119
"zone": true,
124120
"network": true,
@@ -400,8 +396,6 @@ func getJobOptions(ctx context.Context, streaming bool) (*dataflowlib.JobOptions
400396
WorkerHarnessThreads: *workerHarnessThreads,
401397
DiskSizeGb: *diskSizeGb,
402398
DiskType: *diskType,
403-
DiskProvisionedIOPS: *diskProvisionedIOPS,
404-
DiskProvisionedThroughputMibps: *diskProvisionedThroughputMibps,
405399
Algorithm: *autoscalingAlgorithm,
406400
FlexRSGoal: *flexRSGoal,
407401
MachineType: *firstNonEmpty(workerMachineType, machineType),

sdks/go/pkg/beam/runners/dataflow/dataflowlib/job.go

Lines changed: 0 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -58,8 +58,6 @@ type JobOptions struct {
5858
NumWorkers int64
5959
DiskSizeGb int64
6060
DiskType string
61-
DiskProvisionedIOPS int64
62-
DiskProvisionedThroughputMibps int64
6361
MachineType string
6462
Labels map[string]string
6563
ServiceAccountEmail string
@@ -193,8 +191,6 @@ func Translate(ctx context.Context, p *pipepb.Pipeline, opts *JobOptions, worker
193191
},
194192
DiskSizeGb: opts.DiskSizeGb,
195193
DiskType: opts.DiskType,
196-
DiskProvisionedIops: opts.DiskProvisionedIOPS,
197-
DiskProvisionedThroughputMibps: opts.DiskProvisionedThroughputMibps,
198194
IpConfiguration: ipConfiguration,
199195
Kind: "harness",
200196
Packages: packages,

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

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

2929
import ast
3030
import codecs
31-
from functools import partial
3231
import getpass
3332
import hashlib
3433
import io
3534
import json
3635
import logging
3736
import os
3837
import random
39-
import string
40-
41-
from packaging import version
4238
import re
39+
import string
4340
import sys
4441
import time
4542
import traceback
4643
import warnings
4744
from copy import copy
4845
from datetime import datetime
4946
from datetime import timezone
50-
51-
from apitools.base.py import encoding
52-
from apitools.base.py import exceptions
47+
from functools import partial
5348

5449
from apache_beam import version as beam_version
5550
from apache_beam.internal.gcp.auth import get_service_credentials
@@ -73,8 +68,11 @@
7368
from apache_beam.transforms import cy_combiners
7469
from apache_beam.transforms.display import DisplayData
7570
from apache_beam.transforms.environments import is_apache_beam_container
76-
from apache_beam.utils import retry
7771
from apache_beam.utils import proto_utils
72+
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
7876

7977
# Environment version information. It is passed to the service during a
8078
# a job submission and is used by the service to establish what features
@@ -209,8 +207,7 @@ def __init__(
209207
pool.diskProvisionedIops = self.worker_options.disk_provisioned_iops
210208
if self.worker_options.disk_provisioned_throughput_mibps:
211209
pool.diskProvisionedThroughputMibps = (
212-
self.worker_options.disk_provisioned_throughput_mibps
213-
)
210+
self.worker_options.disk_provisioned_throughput_mibps)
214211
if self.worker_options.zone:
215212
pool.zone = self.worker_options.zone
216213
if self.worker_options.network:

0 commit comments

Comments
 (0)