gaborkaszab commented on code in PR #18188:
URL: https://github.com/apache/iceberg/pull/18188#discussion_r4182911966
##########
core/src/main/java/org/apache/iceberg/rest/RESTTableOperations.java:
##########
@@ -165,13 +222,33 @@ public TableMetadata current() {
@Override
public TableMetadata refresh() {
Endpoint.check(endpoints, Endpoint.V1_LOAD_TABLE);
- return updateCurrentMetadata(
+ Map<String, String> responseHeaders =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
+ LoadTableResponse response =
client.get(
path,
readQueryParams,
LoadTableResponse.class,
- readHeaders,
- ErrorHandlers.tableErrorHandler()));
+ this::readHeaders,
+ ErrorHandlers.tableErrorHandler(),
+ responseHeaders::putAll);
+
+ if (response == null) {
+ // metadata is current
+ return current;
+ }
+
+ return updateCurrentMetadata(response,
responseHeaders.get(HttpHeaders.ETAG));
+ }
+
+ private Map<String, String> readHeaders() {
+ Map<String, String> headers = readHeaders.get();
+ if (eTag == null) {
+ return headers;
+ }
+
+ Map<String, String> conditionalHeaders = Maps.newLinkedHashMap(headers);
Review Comment:
nit: I don't think the name `conditionalHeaders` is the best.
`extendedReadHeaders` maybe?
##########
core/src/test/java/org/apache/iceberg/rest/TestFreshnessAwareLoading.java:
##########
@@ -285,6 +285,54 @@ public void notModifiedResponse() {
any());
}
+ @Test
+ public void notModifiedResponseOnRefresh() {
+ restCatalog.createNamespace(TABLE.namespace());
+ restCatalog.createTable(TABLE, SCHEMA);
+ BaseTable table = (BaseTable) restCatalog.loadTable(TABLE);
+ TableMetadata loaded = table.operations().current();
+
+ // the adapter hashes query params into the ETag, so the first refresh
gets a full response
Review Comment:
What are the params used by the catalog adapter that aren't the same used by
ops?
Is this deviation of params general and expected between the catalog client
and ops? Because if we expect this to be the case, then there isn't much point
of constructing ops passing an ETag from `RESTSessionCatalog`, because anyway
the first refresh() is expected to be a full load.
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -4098,6 +4103,209 @@ public void testLoadTableWithSpecialChars() {
assertThat(metadataFileLocation).contains("ns 1 ?=-+/ns 2 ?=-+/table 1
?=-+");
}
+ @Test
+ public void testRefreshSendsIfNoneMatchAndKeepsMetadataOnNotModified() {
Review Comment:
This test seems to cover the a very similar use-case to the one you added in
`TestFreshnessAwareLoading`
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -4098,6 +4103,209 @@ public void testLoadTableWithSpecialChars() {
assertThat(metadataFileLocation).contains("ns 1 ?=-+/ns 2 ?=-+/table 1
?=-+");
}
+ @Test
+ public void testRefreshSendsIfNoneMatchAndKeepsMetadataOnNotModified() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ AtomicReference<Object> lastLoadResponse = new AtomicReference<>();
+ Mockito.doAnswer(
+ invocation -> {
+ Object response = invocation.callRealMethod();
+ lastLoadResponse.set(response);
+ return response;
+ })
+ .when(adapter)
+ .execute(
+ matches(HTTPMethod.GET, RESOURCE_PATHS.table(TABLE)),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String metadataLocation = loaded.metadataFileLocation();
+ // the adapter folds query params into the ETag, so the loadTable and
refresh ETags differ
+ String loadTableETag = ETagProvider.of(metadataLocation,
Map.of("snapshots", "all"));
+ String refreshETag = ETagProvider.of(metadataLocation, Map.of());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshReloadsChangedTableAndAdoptsNewETag() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String loadTableETag =
+ ETagProvider.of(loaded.metadataFileLocation(), Map.of("snapshots",
"all"));
+
+ // another client changes the table
+ backendCatalog
+ .loadTable(TABLE)
+ .updateSchema()
+ .addColumn("extra", Types.LongType.get())
+ .commit();
+
+ TableMetadata refreshed = ops.refresh();
+ assertThat(refreshed).isNotSameAs(loaded);
+
assertThat(refreshed.metadataFileLocation()).isNotEqualTo(loaded.metadataFileLocation());
+ assertThat(refreshed.schema().findField("extra")).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ String refreshETag = ETagProvider.of(refreshed.metadataFileLocation(),
Map.of());
+ assertThat(ops.refresh()).isSameAs(refreshed);
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshUsesETagFromCommit() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
Review Comment:
I think this is always true for these tests.
##########
core/src/main/java/org/apache/iceberg/rest/RESTTableOperations.java:
##########
@@ -222,9 +299,17 @@ public void commit(TableMetadata base, TableMetadata
metadata) {
// the error handler will throw necessary exceptions like
CommitFailedException and
// UnknownCommitStateException
// TODO: ensure that the HTTP client lib passes HTTP client errors to the
error handler
+ Map<String, String> responseHeaders =
Maps.newTreeMap(String.CASE_INSENSITIVE_ORDER);
LoadTableResponse response;
try {
- response = client.post(path, request, LoadTableResponse.class,
mutationHeaders, errorHandler);
+ response =
Review Comment:
The PR title talks about refresh(), but this is commit(). Seems a different
change logically. Or widen the PR title to something like "Freshness-aware
loading in RESTTableOperations".
##########
core/src/test/java/org/apache/iceberg/rest/TestReferencedByQueryParam.java:
##########
@@ -241,12 +243,20 @@ public void refreshKeepsSendingReferencedBy() {
table.refresh();
// refresh() also runs on stale reads and commit retries, and sends no
snapshots parameter
+ String eTag =
Review Comment:
Whenever we decide to send a new query param, this test will fail on the
header verification because of how we calculate the ETag.
Also, the comment above is pretty misleading together with this new code
because the comment says not sending snapshots param, but we anyway use it
right below the comment. (I know why it's needed, but still might be misleading
for the reader)
Do you think you can figure out a more flexible verification? For instance
ETag content is not that relevant for this test, we still make changes so that
we can make the rigid way of our message verification pass. Would be nice to
have a way saying "I'm not interested verifying the outgoing headers" or "I'm
just verifying the name of the outgoing headers but not their values"
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -1439,6 +1442,17 @@ public void testTableAuth(
loaded.refresh(); // refresh to force reload
+ // refresh sends the ETag of the loaded metadata along with the table
headers
+ String eTag =
+ ETagProvider.of(
+ ((BaseTable) loaded).operations().current().metadataFileLocation(),
Review Comment:
Echoing one of my earlier comments: We try to echo here how the server might
populate the ETag. If that changes because, e.g. adding a new param, all the
tests that you modify now will fail. A more flexible verification approach
would help is where we can avoid verified fields that are irrelevant for the
test, E.g. ETag in this case.
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -4098,6 +4103,209 @@ public void testLoadTableWithSpecialChars() {
assertThat(metadataFileLocation).contains("ns 1 ?=-+/ns 2 ?=-+/table 1
?=-+");
}
+ @Test
+ public void testRefreshSendsIfNoneMatchAndKeepsMetadataOnNotModified() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ AtomicReference<Object> lastLoadResponse = new AtomicReference<>();
+ Mockito.doAnswer(
+ invocation -> {
+ Object response = invocation.callRealMethod();
+ lastLoadResponse.set(response);
+ return response;
+ })
+ .when(adapter)
+ .execute(
+ matches(HTTPMethod.GET, RESOURCE_PATHS.table(TABLE)),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String metadataLocation = loaded.metadataFileLocation();
+ // the adapter folds query params into the ETag, so the loadTable and
refresh ETags differ
+ String loadTableETag = ETagProvider.of(metadataLocation,
Map.of("snapshots", "all"));
+ String refreshETag = ETagProvider.of(metadataLocation, Map.of());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshReloadsChangedTableAndAdoptsNewETag() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String loadTableETag =
+ ETagProvider.of(loaded.metadataFileLocation(), Map.of("snapshots",
"all"));
+
+ // another client changes the table
+ backendCatalog
+ .loadTable(TABLE)
+ .updateSchema()
+ .addColumn("extra", Types.LongType.get())
+ .commit();
+
+ TableMetadata refreshed = ops.refresh();
+ assertThat(refreshed).isNotSameAs(loaded);
+
assertThat(refreshed.metadataFileLocation()).isNotEqualTo(loaded.metadataFileLocation());
+ assertThat(refreshed.schema().findField("extra")).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ String refreshETag = ETagProvider.of(refreshed.metadataFileLocation(),
Map.of());
+ assertThat(ops.refresh()).isSameAs(refreshed);
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshUsesETagFromCommit() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ Table table = catalog.loadTable(TABLE);
+ TableOperations ops = ((BaseTable) table).operations();
+
+ table.updateProperties().set("key", "value").commit();
+ TableMetadata committed = ops.current();
+ String commitETag =
Review Comment:
Instead of trying to figure out how the server generated the ETag, we should
capture what the commit returned and verify that it's the same that what the
refresh sent.
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -4098,6 +4103,209 @@ public void testLoadTableWithSpecialChars() {
assertThat(metadataFileLocation).contains("ns 1 ?=-+/ns 2 ?=-+/table 1
?=-+");
}
+ @Test
Review Comment:
All of these tests belong to `TestFreshnessAwareLoading`.
##########
core/src/test/java/org/apache/iceberg/rest/TestRESTCatalog.java:
##########
@@ -4098,6 +4103,209 @@ public void testLoadTableWithSpecialChars() {
assertThat(metadataFileLocation).contains("ns 1 ?=-+/ns 2 ?=-+/table 1
?=-+");
}
+ @Test
+ public void testRefreshSendsIfNoneMatchAndKeepsMetadataOnNotModified() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ AtomicReference<Object> lastLoadResponse = new AtomicReference<>();
+ Mockito.doAnswer(
+ invocation -> {
+ Object response = invocation.callRealMethod();
+ lastLoadResponse.set(response);
+ return response;
+ })
+ .when(adapter)
+ .execute(
+ matches(HTTPMethod.GET, RESOURCE_PATHS.table(TABLE)),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String metadataLocation = loaded.metadataFileLocation();
+ // the adapter folds query params into the ETag, so the loadTable and
refresh ETags differ
+ String loadTableETag = ETagProvider.of(metadataLocation,
Map.of("snapshots", "all"));
+ String refreshETag = ETagProvider.of(metadataLocation, Map.of());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+ assertThat(lastLoadResponse.get()).isNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshReloadsChangedTableAndAdoptsNewETag() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ TableOperations ops = ((BaseTable) catalog.loadTable(TABLE)).operations();
+ TableMetadata loaded = ops.current();
+ String loadTableETag =
+ ETagProvider.of(loaded.metadataFileLocation(), Map.of("snapshots",
"all"));
+
+ // another client changes the table
+ backendCatalog
+ .loadTable(TABLE)
+ .updateSchema()
+ .addColumn("extra", Types.LongType.get())
+ .commit();
+
+ TableMetadata refreshed = ops.refresh();
+ assertThat(refreshed).isNotSameAs(loaded);
+
assertThat(refreshed.metadataFileLocation()).isNotEqualTo(loaded.metadataFileLocation());
+ assertThat(refreshed.schema().findField("extra")).isNotNull();
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, loadTableETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ String refreshETag = ETagProvider.of(refreshed.metadataFileLocation(),
Map.of());
+ assertThat(ops.refresh()).isSameAs(refreshed);
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, refreshETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshUsesETagFromCommit() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ Table table = catalog.loadTable(TABLE);
+ TableOperations ops = ((BaseTable) table).operations();
+
+ table.updateProperties().set("key", "value").commit();
+ TableMetadata committed = ops.current();
+ String commitETag =
+ ETagProvider.of(committed.metadataFileLocation(), Map.of("snapshots",
"all"));
+
+ assertThat(ops.refresh()).isSameAs(committed);
+ verify(adapter)
+ .execute(
+ matches(
+ HTTPMethod.GET,
+ RESOURCE_PATHS.table(TABLE),
+ Map.of(HttpHeaders.IF_NONE_MATCH, commitETag),
+ Map.of()),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+ }
+
+ @Test
+ public void testRefreshWithoutETagFromServer() {
+ RESTCatalogAdapter adapter = Mockito.spy(new
RESTCatalogAdapter(backendCatalog));
+ // the server answers loadTable without an ETag
+ Mockito.doAnswer(
+ invocation -> {
+ Consumer<Map<String, String>> responseHeaders =
invocation.getArgument(3);
+ Consumer<Map<String, String>> withoutETag =
+ headers -> {
+ Map<String, String> filtered = Maps.newHashMap(headers);
+ filtered.remove(HttpHeaders.ETAG);
+ responseHeaders.accept(filtered);
+ };
+ return adapter.execute(
+ invocation.getArgument(0),
+ LoadTableResponse.class,
+ invocation.getArgument(2),
+ withoutETag,
+ ParserContext.builder().build());
+ })
+ .when(adapter)
+ .execute(
+ matches(HTTPMethod.GET, RESOURCE_PATHS.table(TABLE)),
+ eq(LoadTableResponse.class),
+ any(),
+ any());
+
+ RESTCatalog catalog = catalog(adapter);
+ if (requiresNamespaceCreate()) {
+ catalog.createNamespace(TABLE.namespace());
+ }
+
+ catalog.createTable(TABLE, SCHEMA);
+ Table table = catalog.loadTable(TABLE);
+ TableOperations ops = ((BaseTable) table).operations();
+ TableMetadata loaded = ops.current();
+
+ assertThat(ops.refresh()).isSameAs(loaded);
+
+ // the commit response carries an ETag, but the refresh response after it
does not
+ table.updateProperties().set("key", "value").commit();
Review Comment:
I'm not entirely sure what coverage we add with this test. The step with the
commit() probably is not needed. I guess to point is to see what happens when
the server sends no ETag. Then the plan for this test might be:
1) Construct server not to send ETag
2) Create and load table through catalog
3) do a refresh()
4) Verify refresh() didn't send and didn't receive an ETag
5) Verify that the table metadata is as we expect
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]