o-nikolas commented on code in PR #73004:
URL: https://github.com/apache/airflow/pull/73004#discussion_r4097655044
##########
providers/amazon/src/airflow/providers/amazon/aws/log/cloudwatch_task_handler.py:
##########
@@ -148,11 +156,13 @@ def hook(self):
aws_conn_id=conf.get("logging", "remote_log_conn_id"),
region_name=self.region_name
)
- def _build_handler(self) -> watchtower.CloudWatchLogHandler:
+ def _build_handler(self, stream_name: str | None = None) ->
watchtower.CloudWatchLogHandler:
+ if stream_name is None:
+ stream_name = self.log_stream_name
_json_serialize = conf.getimport("aws",
"cloudwatch_task_handler_json_serializer", fallback=None)
return watchtower.CloudWatchLogHandler(
log_group_name=self.log_group,
- log_stream_name=self.log_stream_name,
+ log_stream_name=stream_name,
Review Comment:
`CloudWatchLogHandler.__init__` defaults to `create_log_group=True`, which
runs
`_ensure_log_group()`: a paginated `DescribeLogGroups` call, plus
`CreateLogGroup`
if the group isn't found. On main that happens once per process. With
per-stream
handlers it happens once per trigger, synchronously on the supervisor's log
path,
while holding `_stream_lock`. If it raises, the exception goes straight up
through
the processor into `log.log(...)` in `_process_log_messages_from_subprocess`,
which has no try/except around it.
The handler built when `processors` is set up has already made sure the group
exists, so per-stream handlers can skip the check:
```suggestion
log_stream_name=stream_name,
create_log_group=stream_name is None,
```
---
Drafted-by: Kiro (claude-opus-5.5); reviewed by @o-nikolas before posting
--
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]