jrmccluskey commented on code in PR #40222:
URL: https://github.com/apache/beam/pull/40222#discussion_r4083985327
##########
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:
Let's go with the simpler implementation
##########
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:
done
--
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]