andygrove commented on code in PR #25858:
URL: https://github.com/apache/datafusion/pull/25858#discussion_r4137534498
##########
datafusion/execution/src/cache/cache_manager.rs:
##########
@@ -124,19 +127,36 @@ impl CachedFileMetadata {
/// Check if this cached entry is still valid for the given metadata.
///
/// Returns true if the file size, last modified time, and schema match.
+ /// ETag and version must also match when present in both metadata values.
pub fn is_valid_for(
&self,
current_meta: &ObjectMeta,
current_schema_fingerprint: &Arc<SchemaFingerprint>,
) -> bool {
- self.meta.size == current_meta.size
- && self.meta.last_modified == current_meta.last_modified
+ file_metadata_matches(&self.meta, current_meta)
&& (Arc::ptr_eq(&self.schema_fingerprint,
current_schema_fingerprint)
|| self.schema_fingerprint.as_ref()
== current_schema_fingerprint.as_ref())
}
}
+/// Size and modification time must always match. Some stores or requests omit
+/// object identifiers, so compare each only when both metadata values provide
it.
+fn file_metadata_matches(cached: &ObjectMeta, current: &ObjectMeta) -> bool {
+ cached.size == current.size
+ && cached.last_modified == current.last_modified
+ && cached
+ .e_tag
+ .as_ref()
+ .zip(current.e_tag.as_ref())
+ .is_none_or(|(cached, current)| cached == current)
+ && cached
+ .version
+ .as_ref()
+ .zip(current.version.as_ref())
+ .is_none_or(|(cached, current)| cached == current)
+}
+
Review Comment:
If we compare only ETags when both sides have one, it's worth knowing that
`object_store` passes them through unchanged (as of `object_store` 0.14.2):
- The HTTP store returns the `ETag` header as-is, so weak validators
(`W/"..."`) can show up. Those don't guarantee identical bytes, so an ETag-only
rule should treat a weak ETag as missing.
- For Azure, `list` takes the ETag from the XML `<Etag>` element (unquoted
in `object_store`'s own test fixture, e.g. `0x8D93C7D4629C227`), while `head`
takes the `ETag` response header, which HTTP requires to be quoted. So the same
unchanged blob may not match between list-derived and head-derived
`ObjectMeta`. For example, this happens when a directory table and a
single-file `read_parquet` share the footer cache. That only causes extra cache
misses under either rule, never wrong results, but it might deserve a code
comment.
##########
datafusion/execution/src/cache/mod.rs:
##########
@@ -143,27 +145,71 @@ impl CacheKey for TableScopedPath {
}
}
-/// Each entry is scoped to its use within a specific table so that the cache
-/// can differentiate between identical paths in different tables, and
+/// Identifies an object within its registered object store.
+#[derive(PartialEq, Eq, Hash, Clone, Debug)]
+pub struct ObjectStorePath {
Review Comment:
Naming nit: `ObjectStorePath` is easy to confuse with
`object_store::path::Path`, and `datafusion-cli` already [imports that type
under this exact
alias](https://github.com/apache/datafusion/blob/6ad7b933c4fdcbf94aec6f338fc680618fa38e45/datafusion-cli/src/object_storage/stdin.rs#L32).
What about `StoreScopedPath`, to match `TableScopedPath`? `TableScopedPath`
could then contain one instead of repeating the URL and path fields. Since this
is new public API, it's cheap to change before 56.0.0 ships.
##########
datafusion/catalog-listing/src/table.rs:
##########
@@ -1129,13 +1131,16 @@ impl ListingTable {
store: &Arc<dyn ObjectStore>,
part_file: &PartitionedFile,
) -> datafusion_common::Result<(Arc<Statistics>, Option<LexOrdering>)> {
+ // ListingTable paths share a single object store.
+ let object_store_url = self.table_paths[0].object_store();
Review Comment:
This parses the store URL again (`ObjectStoreUrl::parse`) for every file
during statistics collection. The `store` passed in here is already resolved
from `self.table_paths.first()` by the callers of `collect_files_for_scan`, so
the `ObjectStoreUrl` could be computed once per scan (e.g. in
`collect_files_for_scan`) and passed to this function. That would also remove
the `[0]` index.
##########
datafusion/execution/src/cache/mod.rs:
##########
@@ -143,27 +145,71 @@ impl CacheKey for TableScopedPath {
}
}
-/// Each entry is scoped to its use within a specific table so that the cache
-/// can differentiate between identical paths in different tables, and
+/// Identifies an object within its registered object store.
+#[derive(PartialEq, Eq, Hash, Clone, Debug)]
+pub struct ObjectStorePath {
+ /// URL identifying the registered object store.
+ pub object_store_url: ObjectStoreUrl,
+ /// Location relative to that store.
+ pub path: Path,
+}
+
+impl ObjectStorePath {
+ /// Create a cache key for a path in the given object store.
+ pub fn new(object_store_url: ObjectStoreUrl, path: Path) -> Self {
+ Self {
+ object_store_url,
+ path,
+ }
+ }
+}
+
+impl CacheKey for ObjectStorePath {
+ fn size(&self) -> usize {
+ self.heap_size(&mut DFHeapSizeCtx::default())
+ }
+
+ fn table_ref(&self) -> Option<&TableReference> {
+ None
+ }
+}
+
+impl DFHeapSize for ObjectStorePath {
+ fn heap_size(&self, ctx: &mut DFHeapSizeCtx) -> usize {
+ self.object_store_url.as_str().heap_size(ctx) +
self.path.as_ref().heap_size(ctx)
+ }
+}
+
+impl Display for ObjectStorePath {
Review Comment:
This `Display` impl changes `datafusion-cli`'s `metadata_cache()` output. It
builds the `path` column with
[`path.to_string()`](https://github.com/apache/datafusion/blob/6ad7b933c4fdcbf94aec6f338fc680618fa38e45/datafusion-cli/src/functions.rs#L544),
so paths now include the store URL. `statistics_cache()` and
`list_files_cache()` still print `key.path`, so the same file shows up
differently:
```sql
create external table t stored as parquet location
'/path/to/clickbench_hits_10.parquet';
select count(*) from t;
select 'metadata_cache' as fn, path from metadata_cache()
union all
select 'statistics_cache', path from statistics_cache();
```
```
+------------------+--------------------------------------------+
| fn | path |
+------------------+--------------------------------------------+
| metadata_cache | file:///path/to/clickbench_hits_10.parquet |
| statistics_cache | path/to/clickbench_hits_10.parquet |
+------------------+--------------------------------------------+
```
(output from this PR's `datafusion-cli`; local path replaced with `/path/to`)
The [CLI
docs](https://github.com/apache/datafusion/blob/6ad7b933c4fdcbf94aec6f338fc680618fa38e45/docs/source/user-guide/cli/functions.md?plain=1#L132)
describe this column as "File path relative to the object store / filesystem
root", and their example shows relative paths. The CLI tests don't catch the
change because they only compare [`split_part(path, '/',
-1)`](https://github.com/apache/datafusion/blob/6ad7b933c4fdcbf94aec6f338fc680618fa38e45/datafusion-cli/src/main.rs#L657).
Suggestion: keep `path` store-relative in `metadata_cache()`
(`path.path.to_string()`) and add an `object_store_url` column to all three
functions. Without it, `statistics_cache()` and `list_files_cache()` show
entries from different buckets with the same `path` and `table`, even though
this PR now keeps them apart. It would also be good for the CLI tests to assert
on those columns.
##########
datafusion/datasource/src/url.rs:
##########
@@ -373,6 +375,7 @@ impl ListingTableUrl {
async fn list_with_cache<'b>(
ctx: &'b dyn Session,
store: &'b dyn ObjectStore,
+ object_store_url: &ObjectStoreUrl,
Review Comment:
Minor: the `# Arguments` list in the doc comment above doesn't mention this
new parameter.
--
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]