ahmedabu98 commented on code in PR #40222:
URL: https://github.com/apache/beam/pull/40222#discussion_r4075875770
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java:
##########
@@ -125,6 +129,9 @@ public interface Clock extends Serializable {
*/
public abstract @Nullable Integer getPollingBuckets();
+ @VisibleForTesting
+ public abstract @Nullable Clock getClock();
+
Review Comment:
nit: make this and setClock package private
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java:
##########
@@ -340,7 +414,17 @@ public void processElement(
dynamicDestinations.getTableStringIdentifier(
ValueInSingleWindow.of(element, timestamp, window, paneInfo));
if (tableIdentifier != null && !tableIdentifier.trim().isEmpty()) {
- out.output(tableIdentifier.trim());
+ String canonicalId = tableIdentifier.trim();
Review Comment:
readability nit: I don't think we need the outer `if (tableIdentifier !=
null && !tableIdentifier.trim().isEmpty())` check.
Also let's work directly with `tableIdentifier` instead of trimming it into
a "canonical" version
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java:
##########
@@ -323,10 +342,65 @@ public void populateDisplayData(DisplayData.Builder
builder) {
}
static class ExtractTableIdsDoFn extends DoFn<Row, String> {
+ private static final int DEFAULT_LOCAL_CACHE_MAX_SIZE = 10_000;
+
private final DynamicDestinations dynamicDestinations;
+ private final @Nullable Duration refreshInterval;
+ private final @Nullable Clock clock;
+ private transient @Nullable Ticker ticker;
+ private transient @Nullable Cache<String, Boolean> localTableIdCache;
Review Comment:
Lots of complexity here. I feel like we can simplify with just a simple
`LinkedHashMap<String, Long>` LRU that tracks tableIds and the last emitted
millis. So only emit when `(now - lastEmitted) > refreshInterval` and update
`lastEmitted=now`.
Both ways probably work though so will leave it up to you
##########
sdks/java/io/iceberg/src/test/java/org/apache/beam/sdk/io/iceberg/IcebergIOSideInputTableCacheTest.java:
##########
@@ -607,6 +607,12 @@ public void processElement(@Element Row row,
OutputReceiver<Row> out) {
catalogConfig.catalog().loadTable(IcebergUtils.parseTableIdentifier(tableIdString));
if (table.spec().isUnpartitioned()) {
table.updateSpec().addField("city").commit();
+ // Ensure worker-local table ID cache TTL (interval / 2 = 500ms) has
elapsed
+ try {
+ Thread.sleep(600);
Review Comment:
give this a little extra time to avoid potential flakes?
##########
sdks/java/io/iceberg/src/main/java/org/apache/beam/sdk/io/iceberg/TableMetadataDriver.java:
##########
@@ -340,7 +414,17 @@ public void processElement(
dynamicDestinations.getTableStringIdentifier(
ValueInSingleWindow.of(element, timestamp, window, paneInfo));
if (tableIdentifier != null && !tableIdentifier.trim().isEmpty()) {
- out.output(tableIdentifier.trim());
+ String canonicalId = tableIdentifier.trim();
+ if (localTableIdCache == null && refreshInterval != null) {
+ initCache();
+ }
Review Comment:
readability nit: Do we need this if we're already calling it in setup()?
--
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]