Eliaaazzz commented on code in PR #39461:
URL: https://github.com/apache/beam/pull/39461#discussion_r3672164979


##########
sdks/python/apache_beam/io/fileio.py:
##########
@@ -306,37 +408,81 @@ def __init__(
         '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]:
+    # The Watch transform polls the pattern and emits each file once per key:
+    # the path, joined by the last-modified time when matching updated files.
+    # stop_timestamp bounds the watch to the polls that fall in [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 in PeriodicImpulse. Fall
+        # back to the impulse path — with zero ticks dedup is moot — so the
+        # output stays empty and unbounded, rather than let Watch run its
+        # unconditional first poll.
+        return self._match_all_each_poll(pbegin)
+      termination = _WatchWindowTermination(clock, start_ts.micros, max_polls)
+    if self.match_upd:
+      output_key_fn = _file_path_and_mtime_key
+      output_key_coder = TupleCoder([StrUtf8Coder(), FloatCoder()])

Review Comment:
   You are right, there was a real gap: registry.get_coder does not convert 
typing or native generic annotations, so tuple[str, float] silently fell back 
to pickling, and the explicit coder here was masking exactly that. Watch now 
converts hints with convert_to_beam_type before the registry lookup, the 
explicit settings are gone, and the key functions infer the same StrUtf8Coder 
and TupleCoder as before, with a new watch_test case covering the native 
generic annotation. Also trimmed the docstrings this PR adds and updated the 
stale stacked PR note in the description. This round is one follow up commit on 
top of the reviewed one, commit 4a804a4.



-- 
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]

Reply via email to