tvalentyn commented on code in PR #40297:
URL: https://github.com/apache/beam/pull/40297#discussion_r4125595727
##########
sdks/python/apache_beam/io/gcp/pubsub_integration_test.py:
##########
@@ -138,23 +138,39 @@ def setUp(self):
self.project = self.test_pipeline.get_option('project')
self.uuid = str(uuid.uuid4())
- # Set up PubSub environment.
+ # Set up PubSub environment with retries for transient 504/403 API errors.
from google.cloud import pubsub
self.pub_client = pubsub.PublisherClient()
- self.input_topic = self.pub_client.create_topic(
- name=self.pub_client.topic_path(self.project, INPUT_TOPIC + self.uuid))
- self.output_topic = self.pub_client.create_topic(
- name=self.pub_client.topic_path(self.project, OUTPUT_TOPIC +
self.uuid))
-
self.sub_client = pubsub.SubscriberClient()
- self.input_sub = self.sub_client.create_subscription(
- name=self.sub_client.subscription_path(
- self.project, INPUT_SUB + self.uuid),
- topic=self.input_topic.name)
- self.output_sub = self.sub_client.create_subscription(
- name=self.sub_client.subscription_path(
- self.project, OUTPUT_SUB + self.uuid),
- topic=self.output_topic.name)
+
+ def _retry_pubsub(fn):
+ for attempt in range(4):
+ try:
+ return fn()
+ except Exception:
+ if attempt == 3:
+ raise
+ time.sleep(2**(attempt + 1))
+
+ self.input_topic = _retry_pubsub(
Review Comment:
Let's use existing utils or tenacity for the retry decorator. I'd try
```
from apache_beam.utils import retry
gcp_retry = retry.with_exponential_backoff(
num_retries=3,
initial_delay_secs=2,
retry_filter=retry.retry_on_server_errors_timeout_or_quota_issues_filter
)
self.input_topic = gcp_retry(self.pub_client.create_topic)(
name=self.pub_client.topic_path(self.project, INPUT_TOPIC + self.uuid)
)
...
```
##########
sdks/python/apache_beam/ml/anomaly/transforms_test.py:
##########
@@ -271,9 +277,11 @@ def test_multiple_detectors_without_aggregation(self,
input, expected):
threshold_criterion=FixedThreshold(2),
model_id="zscore_x2"))
- with beam.Pipeline() as p:
+ with TestPipeline(
+ additional_pipeline_args=["--experiments=prism_disable_sdf_split"
Review Comment:
it seems like we are trying to make the test deterministic by enforcing an
order to process the elements.
The problem might be with test assertions if we rely on the order.
cc: @shunping @bvolpato who might have looked into this test before.
+0 on submitting as is since it might help with flakiness.
##########
sdks/python/apache_beam/yaml/integration_tests.py:
##########
@@ -327,6 +327,18 @@ def temp_mongodb_table():
mongo_container.start()
mongo_uri = mongo_container.get_connection_url()
+ # MongoDbContainer's entrypoint restarts mongodb after initial scripts;
wait
+ # until the server is stable, accepting connections and responds to ping.
+ for attempt in range(15):
Review Comment:
nit: we could use tenacity (already a beam dep for retries like this):
```
from tenacity import Retrying, stop_after_attempt, wait_fixed
# MongoDbContainer's entrypoint restarts mongodb after initial scripts; wait
# until the server is stable, accepting connections and responds to ping.
for attempt in Retrying(stop=stop_after_attempt(15), wait=wait_fixed(1)):
with attempt:
mongo_client = mongo_container.get_connection_client()
mongo_client.admin.command('ping')
```
##########
sdks/python/apache_beam/yaml/test_utils/datadog_test_utils.py:
##########
@@ -116,16 +116,14 @@ def temp_fake_datadog_server(expected_records=None):
logging.error(
"Error interacting with temporary fake Datadog server: %s", err)
raise err
- finally:
- if expected_records is not None:
- canonicalize = lambda rec: json.dumps(rec, sort_keys=True)
+ if expected_records is not None:
+ canonicalize = lambda rec: json.dumps(rec, sort_keys=True)
- actual_strs = sorted([canonicalize(r) for r in received])
- expected_strs = sorted([canonicalize(e) for e in expected_records])
+ actual_strs = sorted(set(canonicalize(r) for r in received))
Review Comment:
Is it semantically correct to dedup records by using a set instead of a list?
##########
sdks/python/apache_beam/yaml/integration_tests.py:
##########
@@ -1361,6 +1373,10 @@ def test(self, providers=providers): # default arg to
capture loop value
self.skipTest(
'Runner does not support '
'beam:requirement:pardo:on_window_expiration:v1')
+ if ('VertexAIModelHandlerJSON' in str(exn) and
+ ('PermissionDenied' in str(exn) or 'NotFound' in str(exn) or
+ 'may not exist' in str(exn))):
+ self.skipTest(f'Vertex AI test endpoint unavailable: {exn}')
Review Comment:
> ('PermissionDenied' in str(exn) or 'NotFound' in str(exn)
or
would this silently disable the test due to permission issues?
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]