Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
36 changes: 11 additions & 25 deletions be/src/format_v2/table/lance_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -894,36 +894,22 @@ Status LanceTableReader::_fill_block_from_arrow(LanceBatch* batch, Block* block,
return Status::OK();
}

// The FE sends these already in Lance's own vocabulary, merged from the catalog properties and
// from whatever the namespace vended. Re-encoding them here would drop every option this list did
// not anticipate, so they are handed to lance-c as they arrive.
std::vector<std::string> LanceTableReader::_storage_options(
const TFileScanRangeParams* scan_params) {
if (scan_params == nullptr || !scan_params->__isset.properties) {
if (scan_params == nullptr || !scan_params->__isset.lance_storage_options) {
return {};
}
static constexpr std::array<std::pair<std::string_view, std::string_view>, 5> kStorageKeys = {
{{"AWS_ACCESS_KEY", "aws_access_key_id"},
{"AWS_SECRET_KEY", "aws_secret_access_key"},
{"AWS_TOKEN", "aws_session_token"},
{"AWS_ENDPOINT", "aws_endpoint"},
{"AWS_REGION", "aws_region"}}};
std::vector<std::string> options;
options.reserve(kStorageKeys.size() * 2);
for (const auto& [doris_key, lance_key] : kStorageKeys) {
const auto it = scan_params->properties.find(std::string(doris_key));
if (it != scan_params->properties.end() && !it->second.empty()) {
options.emplace_back(lance_key);
options.emplace_back(it->second);
}
}
const auto endpoint = scan_params->properties.find("AWS_ENDPOINT");
if (endpoint != scan_params->properties.end() && endpoint->second.rfind("http://", 0) == 0) {
options.emplace_back("allow_http");
options.emplace_back("true");
}
const auto path_style = scan_params->properties.find("use_path_style");
if (path_style != scan_params->properties.end() && !path_style->second.empty()) {
const bool use_path_style = path_style->second == "true" || path_style->second == "1";
options.emplace_back("aws_virtual_hosted_style_request");
options.emplace_back(use_path_style ? "false" : "true");
options.reserve(scan_params->lance_storage_options.size() * 2);
for (const auto& [key, value] : scan_params->lance_storage_options) {
if (value.empty()) {
continue;
}
options.emplace_back(key);
options.emplace_back(value);
}
return options;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,10 @@ services:
- ./scripts/lance_rest_server.py:/opt/lance-rest/server.py:ro
environment:
LANCE_REST_BEARER_TOKEN: doris-lance-rest-test-token
LANCE_REST_TABLES_JSON: '{"all_types":"s3://warehouse/lance/all_types.lance"}'
LANCE_REST_TABLES_JSON: '{"all_types":"s3://warehouse/lance/all_types.lance","all_types_unprefixed":"s3://warehouse/lance/all_types.lance"}'
# all_types_unprefixed serves the same dataset but vends its credentials under the
# unprefixed object-store spelling, which is what real namespace servers emit.
LANCE_REST_UNPREFIXED_TABLES_JSON: '["all_types_unprefixed"]'
LANCE_S3_ACCESS_KEY: admin
LANCE_S3_SECRET_KEY: password
LANCE_S3_REGION: us-east-1
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,49 @@ def _load_tables() -> dict[tuple[str, ...], str]:
TABLES = _load_tables()


def _load_unprefixed_tables() -> set[tuple[str, ...]]:
"""Tables whose vended credentials use the unprefixed object-store spelling.

A namespace may spell credentials with any alias Lance accepts, and real servers do use the
unprefixed one, so at least one table has to exercise it.
"""
raw = os.environ.get("LANCE_REST_UNPREFIXED_TABLES_JSON", "[]")
identifiers = json.loads(raw)
if not isinstance(identifiers, list):
raise ValueError("LANCE_REST_UNPREFIXED_TABLES_JSON must be a JSON array")
return {
tuple(part for part in identifier.split(DELIMITER) if part)
for identifier in identifiers
}


UNPREFIXED_TABLES = _load_unprefixed_tables()


def _storage_options(identifier: tuple[str, ...]) -> dict[str, str]:
access_key = os.environ.get("LANCE_S3_ACCESS_KEY", "admin")
secret_key = os.environ.get("LANCE_S3_SECRET_KEY", "password")
region = os.environ.get("LANCE_S3_REGION", "us-east-1")
if identifier in UNPREFIXED_TABLES:
return {
"access_key_id": access_key,
"secret_access_key": secret_key,
"region": region,
"virtual_hosted_style_request": "false",
# Lance refreshes expiring credentials against the namespace itself, so a client has
# to carry this through rather than drop it.
"expires_at_millis": os.environ.get(
"LANCE_S3_EXPIRES_AT_MILLIS", "4102444800000"
),
}
return {
"aws_access_key_id": access_key,
"aws_secret_access_key": secret_key,
"aws_region": region,
"aws_virtual_hosted_style_request": "false",
}


def _decode_identifier(identifier: str) -> tuple[str, ...]:
identifier = unquote(identifier)
if identifier == DELIMITER:
Expand Down Expand Up @@ -141,16 +184,7 @@ def do_POST(self) -> None:
"namespace": list(identifier[:-1]),
"location": table_uri,
"table_uri": table_uri,
"storage_options": {
"aws_access_key_id": os.environ.get(
"LANCE_S3_ACCESS_KEY", "admin"
),
"aws_secret_access_key": os.environ.get(
"LANCE_S3_SECRET_KEY", "password"
),
"aws_region": os.environ.get("LANCE_S3_REGION", "us-east-1"),
"aws_virtual_hosted_style_request": "false",
},
"storage_options": _storage_options(identifier),
"managed_versioning": False,
"is_only_declared": False,
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,6 @@
import java.util.ArrayDeque;
import java.util.ArrayList;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
import java.util.LinkedHashSet;
import java.util.List;
Expand Down Expand Up @@ -79,8 +78,7 @@ public class LanceExternalCatalog extends ExternalCatalog {
private transient List<String> parentNamespace = Collections.emptyList();
private transient String catalogType;
private transient String rootDatabase;
private transient Map<String, String> javaStorageOptions = Collections.emptyMap();
private transient Map<String, String> backendStorageOptions = Collections.emptyMap();
private transient Map<String, String> lanceStorageOptions = Collections.emptyMap();
private transient Object namespaceLock = new Object();

public LanceExternalCatalog(long catalogId, String name, String resource, Map<String, String> props,
Expand All @@ -98,11 +96,11 @@ protected void initLocalObjectsImpl() {
rootDatabase = properties.getRootDatabase();
parentNamespace = LanceNamespaceName.parseParentNamespace(
properties.getNamespaceParent(), properties.getNamespaceDelimiter());
backendStorageOptions = catalogProperty.getBackendStorageProperties();
javaStorageOptions = LanceStorageOptions.forJavaSdk(backendStorageOptions);
lanceStorageOptions = LanceStorageOptions.toLanceOptions(
catalogProperty.getBackendStorageProperties());

allocator = new RootAllocator(ALLOCATOR_LIMIT);
namespace = properties.createNamespace(allocator, javaStorageOptions);
namespace = properties.createNamespace(allocator, lanceStorageOptions);
} catch (Exception e) {
closeLanceObjects();
throw new RuntimeException("Failed to initialize Lance catalog '" + getName()
Expand All @@ -120,7 +118,7 @@ public void checkWhenCreating() throws DdlException {
}

AbstractLanceProperties properties = getLanceProperties();
Map<String, String> storageOptions = LanceStorageOptions.forJavaSdk(
Map<String, String> storageOptions = LanceStorageOptions.toLanceOptions(
catalogProperty.getBackendStorageProperties());
List<String> parent = LanceNamespaceName.parseParentNamespace(
properties.getNamespaceParent(), properties.getNamespaceDelimiter());
Expand Down Expand Up @@ -285,12 +283,10 @@ public LanceTableMetadata loadTableMetadata(String dbName, String tableName,
throw new RuntimeException("Lance namespace returned no table URI for " + dbName + "." + tableName);
}

Map<String, String> storageOptions = new HashMap<>(javaStorageOptions);
if (table.getStorageOptions() != null) {
storageOptions.putAll(table.getStorageOptions());
}
Map<String, String> tableBackendStorageOptions = LanceStorageOptions.forBackend(
backendStorageOptions, table.getStorageOptions());
// One option map serves both readers: the FE opens the dataset through the Lance Java SDK
// and the BE through lance-c, so neither can end up with credentials the other lacks.
Map<String, String> storageOptions = LanceStorageOptions.mergeVended(
lanceStorageOptions, table.getStorageOptions());
try {
if (tableSnapshot.isPresent()) {
TableSnapshot snapshot = tableSnapshot.get();
Expand All @@ -306,11 +302,9 @@ public LanceTableMetadata loadTableMetadata(String dbName, String tableName,
version = LanceSnapshotResolver.getVersionAtOrBefore(
datasetUri, storageOptions, timestamp, allocator);
}
return LanceMetadataLoader.loadVersion(datasetUri, storageOptions,
tableBackendStorageOptions, version, allocator);
return LanceMetadataLoader.loadVersion(datasetUri, storageOptions, version, allocator);
}
return LanceMetadataLoader.loadLatest(datasetUri, storageOptions,
tableBackendStorageOptions, allocator);
return LanceMetadataLoader.loadLatest(datasetUri, storageOptions, allocator);
} catch (Exception e) {
throw new RuntimeException("Failed to load Lance table metadata for " + dbName + "." + tableName
+ ": " + sanitizedRootCauseMessage(e), safeCause(e));
Expand Down Expand Up @@ -348,11 +342,6 @@ private List<String> buildFullNamespace(List<String> relativeNamespace) {
return result;
}

public Map<String, String> getBackendStorageOptions() {
makeSureInitialized();
return backendStorageOptions;
}

public String getLanceCatalogType() {
makeSureInitialized();
return catalogType;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,10 @@ private LanceMetadataLoader() {
* Lance dataset through an S3 TVF.
*/
public static LanceTableMetadata loadLatestForTvf(
String datasetUri, Map<String, String> backendStorageOptions)
String datasetUri, Map<String, String> backendProperties)
throws Exception {
try (BufferAllocator allocator = new RootAllocator(ALLOCATOR_LIMIT)) {
return loadLatest(datasetUri, LanceStorageOptions.forJavaSdk(backendStorageOptions),
backendStorageOptions, allocator);
return loadLatest(datasetUri, LanceStorageOptions.toLanceOptions(backendProperties), allocator);
}
}

Expand All @@ -57,10 +56,9 @@ public static LanceTableMetadata loadLatestForTvf(
* time-travel version is requested. Schema, version, and fragments are read from the same
* opened dataset snapshot.
*/
public static LanceTableMetadata loadLatest(String datasetUri, Map<String, String> javaStorageOptions,
Map<String, String> backendStorageOptions, BufferAllocator allocator) throws Exception {
return loadInternal(
datasetUri, javaStorageOptions, backendStorageOptions, OptionalLong.empty(), allocator);
public static LanceTableMetadata loadLatest(String datasetUri, Map<String, String> lanceStorageOptions,
BufferAllocator allocator) throws Exception {
return loadInternal(datasetUri, lanceStorageOptions, OptionalLong.empty(), allocator);
}

/**
Expand All @@ -70,18 +68,16 @@ public static LanceTableMetadata loadLatest(String datasetUri, Map<String, Strin
* {@link LanceExternalCatalog#loadTableMetadata(String, String, java.util.Optional)} for both
* {@code FOR VERSION AS OF} and the version resolved from {@code FOR TIME AS OF}.
*/
public static LanceTableMetadata loadVersion(String datasetUri, Map<String, String> javaStorageOptions,
Map<String, String> backendStorageOptions, long version, BufferAllocator allocator) throws Exception {
return loadInternal(
datasetUri, javaStorageOptions, backendStorageOptions, OptionalLong.of(version), allocator);
public static LanceTableMetadata loadVersion(String datasetUri, Map<String, String> lanceStorageOptions,
long version, BufferAllocator allocator) throws Exception {
return loadInternal(datasetUri, lanceStorageOptions, OptionalLong.of(version), allocator);
}

/** Shared implementation for the latest-version and explicit-version public entry points. */
private static LanceTableMetadata loadInternal(String datasetUri, Map<String, String> javaStorageOptions,
Map<String, String> backendStorageOptions, OptionalLong version,
BufferAllocator allocator) throws Exception {
private static LanceTableMetadata loadInternal(String datasetUri, Map<String, String> lanceStorageOptions,
OptionalLong version, BufferAllocator allocator) throws Exception {
try (Dataset dataset = Dataset.open().allocator(allocator).uri(datasetUri)
.readOptions(LanceReadOptions.build(javaStorageOptions, version)).build()) {
.readOptions(LanceReadOptions.build(lanceStorageOptions, version)).build()) {
long resolvedVersion = dataset.version();
List<LanceTableMetadata.LanceFragmentInfo> fragments = new ArrayList<>();
for (Fragment fragment : dataset.getFragments()) {
Expand All @@ -90,7 +86,7 @@ private static LanceTableMetadata loadInternal(String datasetUri, Map<String, St
fragment.metadata().getPhysicalRows()));
}
return new LanceTableMetadata(datasetUri, resolvedVersion, dataset.getSchema(), fragments,
backendStorageOptions);
lanceStorageOptions);
}
}
}
Loading
Loading