viukpe opened a new pull request, #74433:
URL: https://github.com/apache/airflow/pull/74433
### What happens
When you run a Beam pipeline operator — BeamRunPythonPipelineOperator,
BeamRunJavaPipelineOperator, BeamRunGoPipelineOperator, or a Google Dataflow
operator that launches through the shared Beam hook — the task launches the
pipeline as a subprocess and streams its stdout/stderr back to the task log.
That streaming is done in run_beam_command, which waits for output using
Python's stdlib select.select():
reads = [proc.stderr, proc.stdout]
readable_fds, _, _ = select.select(reads, [], [], 5)
select.select() is backed by a fixed-size fd_set bitmap of FD_SETSIZE (1024)
bits. It physically cannot represent a file descriptor whose number is >=
1024,
and raises "ValueError: filedescriptor out of range in select()" when asked
to.
This ceiling is compile-time; it is unrelated to RLIMIT_NOFILE, so raising
ulimit -n does not help.
So whether this fires depends only on the fd NUMBER assigned to the
subprocess
pipes. File descriptors are handed out lowest-free-first, so when the task
process already holds many descriptors open, the Beam subprocess's
stdout/stderr
pipes are assigned numbers >= 1024 — and the pipeline run aborts with the
error
above, even though nothing is actually wrong with the pipeline.
This is the exact same root cause and class of fix already accepted for the
SSH
provider in #74211 (hook) and #74299 (tunnel), and it traces back to the same
FD_SETSIZE issue I originally reported in #74205. The synchronous Beam path
is
affected; the async path (run_beam_command_async, asyncio-based) is not.
### Fix
Replace the select.select() loop with selectors.DefaultSelector (epoll/poll),
which has no FD_SETSIZE ceiling. Behaviour is otherwise unchanged.
### Reproduction / validation
- Reproduced on the real run_beam_command with a real subprocess while the
process held 1,100+ open fds: the old select() code raises
"filedescriptor out of range in select()"; the fix captures the subprocess
output correctly.
- With few fds open (< 1024) both behave identically — the only variable that
flips pass/fail is the fd number.
- Updated the existing logging test to patch selectors; added a regression
test
that runs a real subprocess under a high fd count (fails on old select(),
passes with selectors). ruff check + format clean.
related: #74205
--
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]