Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
17 commits
Select commit Hold shift + click to select a range
3760698
[Python] Refactor MatchContinuously onto the Watch transform
Eliaaazzz Jul 23, 2026
4a804a4
Address review: close the coder inference gap, drop explicit key coders
Eliaaazzz Jul 29, 2026
4587416
Annotate poll output type, cover typing.Tuple inference, sort test im…
Eliaaazzz Jul 29, 2026
b8b96ae
Split the coder inference fix into #39547
Eliaaazzz Jul 29, 2026
ba21f51
Merge remote-tracking branch 'upstream/master' into matchcontinuously…
Eliaaazzz Aug 10, 2026
8a132af
Add a timestamp_cursor option to MatchContinuously
Eliaaazzz Aug 10, 2026
462d2f5
Keep the Watch restriction's microsecond precision
Eliaaazzz Aug 12, 2026
afa436c
Address review on the MatchContinuously timestamp cursor
Eliaaazzz Aug 12, 2026
f68f51e
Rewrap the timestamp_cursor doc paragraph
Eliaaazzz Aug 12, 2026
9df6713
Retire keys with the Watch cursor instead of replacing them
Eliaaazzz Aug 13, 2026
9ce7c97
Let the MatchContinuously cursor compose with the match key
Eliaaazzz Aug 13, 2026
0526f7d
Release the mtime watermark when a poll finds nothing newer
Eliaaazzz Aug 13, 2026
f184d2f
Revert "Release the mtime watermark when a poll finds nothing newer"
Eliaaazzz Aug 13, 2026
626ee74
Release the mtime watermark when a poll finds nothing newer
Eliaaazzz Aug 13, 2026
7669d4a
[Python] Key the MatchContinuously cursor on the path and mtime
Eliaaazzz Aug 14, 2026
1a56c7d
[Python] Imply match_updated_files under the timestamp cursor
Eliaaazzz Aug 17, 2026
acbf862
[Python] Cut the timestamp cursor docs down to the contract
Eliaaazzz Aug 17, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
305 changes: 219 additions & 86 deletions sdks/python/apache_beam/io/fileio.py
Original file line number Diff line number Diff line change
Expand Up @@ -93,28 +93,35 @@
import random
import uuid
from collections import namedtuple
from functools import partial
from typing import Any
from typing import BinaryIO # pylint: disable=unused-import
from typing import Callable
from typing import Iterable
from typing import Optional
from typing import Union

import apache_beam as beam
from apache_beam.coders.coders import VarIntCoder
from apache_beam.io import filesystem
from apache_beam.io import filesystems
from apache_beam.io.filesystem import BeamIOError
from apache_beam.io.filesystem import CompressionTypes
from apache_beam.io.watch import PollFn
from apache_beam.io.watch import PollResult
from apache_beam.io.watch import TerminationCondition
from apache_beam.io.watch import Watch
from apache_beam.io.watch import never
from apache_beam.options.pipeline_options import GoogleCloudOptions
from apache_beam.options.value_provider import StaticValueProvider
from apache_beam.options.value_provider import ValueProvider
from apache_beam.transforms.periodicsequence import PeriodicImpulse
from apache_beam.transforms.userstate import CombiningValueStateSpec
from apache_beam.transforms.window import BoundedWindow
from apache_beam.transforms.window import FixedWindows
from apache_beam.transforms.window import GlobalWindow
from apache_beam.transforms.window import IntervalWindow
from apache_beam.transforms.window import TimestampedValue
from apache_beam.utils.timestamp import MAX_TIMESTAMP
from apache_beam.utils.timestamp import Duration
from apache_beam.utils.timestamp import Timestamp

__all__ = [
Expand Down Expand Up @@ -251,6 +258,115 @@ def process(
yield ReadableFile(metadata, self._compression)


class _PollClock(object):
"""Shares one clock reading per poll round, so the start gate and the poll
budget judge the ``start_timestamp`` boundary consistently."""
def __init__(self):
self.last_poll_micros: Optional[int] = None


class _WatchWindowTermination(TerminationCondition):
"""Stops after the polls that fall in the ``[start, stop)`` window.

``max_polls`` is the ``PeriodicImpulse`` tick count
``ceil((stop - start) / interval)``; polls before ``start`` are waiting
rounds and do not consume the budget.
"""
def __init__(self, clock: _PollClock, start_micros: int, max_polls: int):
self._clock = clock
self._start_micros = start_micros
self._max_polls = max_polls

def for_new_input(self, now, element):
return 0

def on_poll_complete(self, state):
poll_micros = self._clock.last_poll_micros
if poll_micros is not None and poll_micros >= self._start_micros:
return state + 1
return state

def can_stop_polling(self, now, state):
return state >= self._max_polls

def state_coder(self):
return VarIntCoder()


def _ensure_mtime(metadata: filesystem.FileMetadata) -> float:
# A missing (zero) timestamp is rejected because every file would then carry
# the same one, and updates could never be told apart.
if not metadata.last_updated_in_seconds:
raise BeamIOError(
'MatchContinuously deduplicates by last-modified time, but %s reports '
'none.' % metadata.path)
return metadata.last_updated_in_seconds


def _file_path_key(metadata: filesystem.FileMetadata) -> str:
return metadata.path


def _file_path_and_mtime_key(
metadata: filesystem.FileMetadata) -> tuple[str, float]:
# Keying on the last-modified time makes a changed file look new again.
return metadata.path, _ensure_mtime(metadata)


class _MatchContinuouslyPollFn(PollFn):
"""Polls a file pattern, honoring empty-match rules.

A poll before ``start_timestamp`` emits nothing. Matches carry the poll time
as their event time, or their last-modified time under ``mtime_timestamps``,
where the watermark trails the newest last-modified time for as long as
polls keep turning up newer ones.
"""
def __init__(
self,
empty_match_treatment,
start_timestamp,
clock=None,
mtime_timestamps=False):
self._empty_match_treatment = empty_match_treatment
self._start_micros = Timestamp.of(start_timestamp).micros
self._clock = clock if clock is not None else _PollClock()
self._mtime_timestamps = mtime_timestamps
# Greatest last-modified time handed out so far, to tell a poll that found
# something newer from one that only re-listed what was already there.
self._newest_mtime = None # type: Optional[Timestamp]

def __call__(self, file_pattern: str) -> PollResult[filesystem.FileMetadata]:
now = Timestamp.now()
self._clock.last_poll_micros = now.micros
if now.micros < self._start_micros:
return PollResult.incomplete(())
match_result = filesystems.FileSystems.match([file_pattern])[0]
if (not match_result.metadata_list and
not EmptyMatchTreatment.allow_empty_match(file_pattern,
self._empty_match_treatment)):
raise BeamIOError(
'Empty match for pattern %s. Disallowed.' % file_pattern)
if not self._mtime_timestamps:
return PollResult.incomplete(
match_result.metadata_list, timestamp=now).with_watermark(now)
outputs = [
TimestampedValue(metadata, Timestamp.of(_ensure_mtime(metadata)))
for metadata in match_result.metadata_list
]
# A poll that turned up a newer last-modified time has just read the
# filesystem clock, so the watermark stops there, capped at the poll time.
# A poll that found nothing newer takes the poll time, so a quiet
# directory does not stall event-time windows.
newest = max((output.timestamp for output in outputs), default=None)
if newest is not None and (self._newest_mtime is None or
newest > self._newest_mtime):
self._newest_mtime = newest
watermark = min(newest, now)
else:
watermark = now
return PollResult.incomplete(outputs).with_watermark(watermark)


class MatchContinuously(beam.PTransform):
"""Checks for new files for a given pattern every interval.

Expand All @@ -260,12 +376,16 @@ class MatchContinuously(beam.PTransform):
MatchContinuously is experimental. No backwards-compatibility
guarantees.

Matching continuously scales poorly, as it is stateful, and requires storing
file ids in memory. In addition, because it is memory-only, if a pipeline is
restarted, already processed files will be reprocessed. Consider an alternate
technique, such as Pub/Sub Notifications
(https://cloud.google.com/storage/docs/pubsub-notifications)
when using GCS if possible.
Deduplication state is checkpointed, so a runner with checkpointing enabled
restores it after a restart and does not reprocess files. That state grows
with the number of files matched, unless ``timestamp_cursor`` bounds it. For
a growing directory on GCS, consider an alternate technique such as Pub/Sub
Notifications (https://cloud.google.com/storage/docs/pubsub-notifications).

A match carries the poll time as its event time, and the watermark follows
the poll time. Under ``timestamp_cursor`` a match carries its last-modified
time instead, and the watermark holds at the newest one matched, capped at
the poll time, until a poll turns up nothing newer and releases it.
"""
def __init__(
self,
Expand All @@ -276,7 +396,8 @@ def __init__(
stop_timestamp=MAX_TIMESTAMP,
match_updated_files=False,
apply_windowing=False,
empty_match_treatment=EmptyMatchTreatment.ALLOW):
empty_match_treatment=EmptyMatchTreatment.ALLOW,
timestamp_cursor=False):
"""Initializes a MatchContinuously transform.

Args:
Expand All @@ -289,6 +410,12 @@ def __init__(
file with timestamp changes.
apply_windowing: Whether each element should be assigned to
individual window. If false, all elements will reside in global window.
timestamp_cursor: (When match_updated_files and has_deduplication are set
to True) bound the deduplication state by last-modified time. By
default, all file modification history is tracked. If set to true, file
modification history prior to the max(mtime of last poll result) are
dropped, for better performance. A file that appears with an older
last-modified time is then taken as already seen and skipped.
"""

self.file_pattern = file_pattern
Expand All @@ -299,44 +426,97 @@ def __init__(
self.match_upd = match_updated_files
self.apply_windowing = apply_windowing
self.empty_match_treatment = empty_match_treatment
_LOGGER.warning(
'Matching Continuously is stateful, and can scale poorly. '
'Consider using Pub/Sub Notifications '
'(https://cloud.google.com/storage/docs/pubsub-notifications) '
'if possible')
self.timestamp_cursor = timestamp_cursor
if timestamp_cursor:
if not has_deduplication:
raise ValueError(
'MatchContinuously(timestamp_cursor=True) deduplicates, so it '
'requires has_deduplication=True.')
if not match_updated_files:
_LOGGER.warning(
'MatchContinuously(timestamp_cursor=True) implies '
'match_updated_files=True.')
self.match_upd = True
else:
_LOGGER.warning(
'Matching Continuously is stateful, and can scale poorly. '
'Consider using Pub/Sub Notifications '
'(https://cloud.google.com/storage/docs/pubsub-notifications) '
'if possible')

def expand(self, pbegin) -> beam.PCollection[filesystem.FileMetadata]:
# invoke periodic impulse
impulse = pbegin | PeriodicImpulse(
start_timestamp=self.start_ts,
stop_timestamp=self.stop_ts,
fire_interval=self.interval)

# match file pattern periodically
file_pattern = self.file_pattern
match_files = (
impulse
| 'GetFilePattern' >> beam.Map(lambda x: file_pattern)
| MatchAll(self.empty_match_treatment))

# apply deduplication strategy if required
if Duration.of(self.interval).micros <= 0:
raise ValueError('MatchContinuously interval must be positive.')
if self.has_deduplication:
# Making a Key Value so each file has its own state.
match_files = match_files | 'ToKV' >> beam.Map(lambda x: (x.path, x))
if self.match_upd:
match_files = match_files | 'RemoveOldAlreadyRead' >> beam.ParDo(
_RemoveOldDuplicates())
else:
match_files = match_files | 'RemoveAlreadyRead' >> beam.ParDo(
_RemoveDuplicates())

# apply windowing if required. Apply at last because deduplication relies on
# the global window.
match_files = self._match_deduplicated(pbegin)
else:
match_files = self._match_all_each_poll(pbegin)

# Apply windowing last because dedup relies on the global window.
if self.apply_windowing:
match_files = match_files | beam.WindowInto(FixedWindows(self.interval))

return match_files

def _match_deduplicated(self,
pbegin) -> beam.PCollection[filesystem.FileMetadata]:
# Watch emits each file once per dedup key: the path, joined by the mtime
# when matching updated files. stop_timestamp bounds the polls to
# [start, stop).
clock = _PollClock()
if self.stop_ts == MAX_TIMESTAMP:
termination = never()
else:
start_ts = Timestamp.of(self.start_ts)
stop_ts = Timestamp.of(self.stop_ts)
if stop_ts < start_ts:
raise ValueError(
'MatchContinuously stop_timestamp %s precedes start_timestamp %s' %
(stop_ts, start_ts))
interval_micros = Duration.of(self.interval).micros
span_micros = (stop_ts - start_ts).micros
# Ceiling division reproduces PeriodicImpulse's tick count; the window
# upper bound is exclusive.
max_polls = -(-span_micros // interval_micros)
if max_polls == 0:
# An empty [start, stop) window never ticks; the impulse path keeps
# the output empty without Watch's unconditional first poll.
return self._match_all_each_poll(pbegin)
termination = _WatchWindowTermination(clock, start_ts.micros, max_polls)
poll_fn = _MatchContinuouslyPollFn(
self.empty_match_treatment,
self.start_ts,
clock,
mtime_timestamps=self.timestamp_cursor)
# The key coder is inferred from the key function's return annotation.
watch = Watch(
poll_fn,
poll_interval=self.interval,
termination=termination,
output_key_fn=(
_file_path_and_mtime_key if self.match_upd else _file_path_key),
timestamp_cursor=self.timestamp_cursor)
# Watch emits (pattern, file) pairs; keep the FileMetadata output type so
# downstream transforms stay typed instead of falling back to Any.
return (
pbegin
| 'Impulse' >> beam.Create([self.file_pattern])
| 'Watch' >> watch
| 'DropPattern' >> beam.Map(lambda kv: kv[1]).with_output_types(
filesystem.FileMetadata))

def _match_all_each_poll(self,
pbegin) -> beam.PCollection[filesystem.FileMetadata]:
# No deduplication: re-emit every match on each poll.
return (
pbegin
| PeriodicImpulse(
start_timestamp=self.start_ts,
stop_timestamp=self.stop_ts,
fire_interval=self.interval)
| 'GetFilePattern' >> beam.Map(lambda x: self.file_pattern)
| MatchAll(self.empty_match_treatment))


class ReadMatches(beam.PTransform):
"""Converts each result of MatchFiles() or MatchAll() to a ReadableFile.
Expand Down Expand Up @@ -892,50 +1072,3 @@ def finish_bundle(self):
timestamp=key[1].start,
windows=[key[1]] # TODO(pabloem) HOW DO WE GET THE PANE
))


class _RemoveDuplicates(beam.DoFn):
"""Internal DoFn that filters out filenames already seen (even though the file
has updated)."""
COUNT_STATE = CombiningValueStateSpec('count', combine_fn=sum)

def process(
self,
element: tuple[str, filesystem.FileMetadata],
count_state=beam.DoFn.StateParam(COUNT_STATE)
) -> Iterable[filesystem.FileMetadata]:

path = element[0]
file_metadata = element[1]
counter = count_state.read()

if counter == 0:
count_state.add(1)
_LOGGER.debug('Generated entry for file %s', path)
yield file_metadata
else:
_LOGGER.debug('File %s was already read, seen %d times', path, counter)


class _RemoveOldDuplicates(beam.DoFn):
"""Internal DoFn that filters out filenames already seen and timestamp
unchanged."""
TIME_STATE = CombiningValueStateSpec(
'count', combine_fn=partial(max, default=0.0))

def process(
self,
element: tuple[str, filesystem.FileMetadata],
time_state=beam.DoFn.StateParam(TIME_STATE)
) -> Iterable[filesystem.FileMetadata]:
path = element[0]
file_metadata = element[1]
new_ts = file_metadata.last_updated_in_seconds
old_ts = time_state.read()

if old_ts < new_ts:
time_state.add(new_ts)
_LOGGER.debug('Generated entry for file %s', path)
yield file_metadata
else:
_LOGGER.debug('File %s was already read', path)
Loading
Loading