Ce guide explique comment lire et écrire des tables Lance gérées par DLF à l'aide du moteur DataFrame Daft. Daft propose une API DataFrame lazy bien adaptée aux scénarios de filtrage de requêtes et de calcul par lots.
Pour lire et écrire directement des tables Lance avec PyLance, consultez Utilisation de DLF Lance en Python.
Terminologie
|
Composant |
Rôle |
|
DLF |
Service de catalogue qui gère les métadonnées de bases de données et de tables, stocke les chemins d'accès aux tables Lance et délivre des identifiants OSS temporaires |
|
Lance/PyLance |
Format de données et implémentation de lecture/écriture de bas niveau responsable des opérations d'E/S réelles sur les ensembles de données Lance hébergés sur OSS |
|
Daft |
Moteur de calcul DataFrame fournissant les interfaces |
|
|
Connecteur qui récupère les chemins d'accès aux tables et les identifiants OSS temporaires depuis DLF pour utilisation par Daft. N'expose que les tables dont le |
Architecture
User code
→ lance_namespace.connect("dlf", CONFIG) # Connect to DLF catalog
→ DLF returns Lance table path + temporary OSS credentials
→ apply_oss_environment(...) # Set OSS_* environment variables
→ daft.read_lance("oss://...") # Read data
→ df.write_lance("oss://...", mode="append") # Write data
Correspondance des concepts :
Base de données DLF → namespace Lance
Table DLF → table Lance
Prérequis
Installer les dépendances
python3 -m pip install lance-dlf daft
Le package lance-dlf installe automatiquement lance_namespace, pyarrow et les autres dépendances requises.
Configuration minimale
CONFIG = {
"uri": "http://<dlf-endpoint>",
"warehouse": "<warehouse>",
"token.provider": "dlf",
"dlf.region": "<region>",
"dlf.access-key-id": "<access-key-id>",
"dlf.access-key-secret": "<access-key-secret>",
"dlf.oss-endpoint": "<oss-endpoint>",
}
Se connecter à DLF
L'importation de lance_dlf enregistre automatiquement le namespace dlf :
import lance_namespace
import lance_dlf # noqa: F401
ns = lance_namespace.connect("dlf", CONFIG)
print(ns.namespace_id())
Lire une table existante
Étape 1 : Obtenir le chemin d'accès à la table et les identifiants
Avant de lire ou d'écrire une table avec Daft, appelez describe_table() pour récupérer le chemin d'accès à la table et les identifiants OSS temporaires depuis DLF.
from lance_namespace import DescribeTableRequest
DATABASE = "<database>"
TABLE = "<table>"
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
print(desc.location)
print(sorted((desc.storage_options or {}).keys()))
L'objet desc contient deux champs clés :
desc.location— le chemin de stockage de la table Lance au formatoss://bucket/path/to/tabledesc.storage_options— un dictionnaire contenant les identifiants OSS temporaires
Étape 2 : Définir les variables d'environnement des identifiants OSS
Daft utilise PyLance en arrière-plan pour accéder à OSS. Définissez une fonction d'assistance apply_oss_environment qui mappe les identifiants délivrés par DLF vers les variables d'environnement OSS_*, puis appelez-la :
import os
def apply_oss_environment(storage_options: dict) -> None:
os.environ["OSS_ENDPOINT"] = storage_options["oss_endpoint"]
os.environ["OSS_ACCESS_KEY_ID"] = storage_options["oss_access_key_id"]
os.environ["OSS_ACCESS_KEY_SECRET"] = storage_options["oss_secret_access_key"]
if storage_options.get("oss_security_token"):
os.environ["OSS_SECURITY_TOKEN"] = storage_options["oss_security_token"]
if storage_options.get("oss_region"):
os.environ["OSS_REGION"] = storage_options["oss_region"]
# Apply credentials
apply_oss_environment(desc.storage_options or {})
Étape 3 : Lire les données de la table avec Daft
import daft
df = daft.read_lance(desc.location)
df.show()
Écrire des données
Ajouter des données à une table existante
Après avoir terminé l'Étape 1 : Obtenir le chemin d'accès à la table et les identifiants et l'Étape 2 : Définir les variables d'environnement des identifiants OSS ci-dessus, ajoutez des données avec mode="append" :
# Prerequisites: connect to Catalog, get table path/credentials, set OSS env vars
desc = ns.describe_table(DescribeTableRequest(id=[DATABASE, TABLE]))
apply_oss_environment(desc.storage_options or {})
# Append data
append_df = daft.from_pydict({
"f0": [204],
"f1": ["daft-d"],
})
append_df.write_lance(desc.location, mode="append")
# Verify the write
df2 = daft.read_lance(desc.location)
df2.show()
Créer une nouvelle table et y écrire des données
Les nouvelles tables doivent d'abord être créées avec ns.create_table() (qui écrit également le premier lot de données). Utilisez Daft pour les lectures et écritures ultérieures.
# Prerequisites: connect to Catalog, get table path/credentials, set OSS env vars
from datetime import datetime
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
# Serialize an Arrow table to IPC bytes
def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, table.schema) as writer:
writer.write_table(table)
return sink.getvalue().to_pybytes()
# Create the table with initial data
table_name = "test_lance_daft_" + datetime.now().strftime("%Y%m%d_%H%M%S")
table_id = [DATABASE, table_name]
rows = {
"f0": [201, 202, 203],
"f1": ["daft-a", "daft-b", "daft-c"],
}
arrow_table = pa.table(rows)
create_response = ns.create_table(
CreateTableRequest(id=table_id),
arrow_table_to_ipc_bytes(arrow_table),
)
print(create_response.location)
Après la création, récupérez les identifiants avec describe_table et utilisez Daft pour lire et écrire :
# Get credentials and set environment variables
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# Read and verify
df = daft.read_lance(desc.location)
df.show()
# Append data with Daft
append_rows = {
"f0": [204],
"f1": ["daft-d"],
}
append_df = daft.from_pydict(append_rows)
meta = append_df.write_lance(desc.location, mode="append")
meta.show()
# Read again to confirm
appended_df = daft.read_lance(desc.location)
appended_df.show()
Sortie attendue :
[
{"f0": 201, "f1": "daft-a"},
{"f0": 202, "f1": "daft-b"},
{"f0": 203, "f1": "daft-c"},
{"f0": 204, "f1": "daft-d"}
]
Exemple complet
Le script suivant illustre le flux de travail de bout en bout : créer une nouvelle table → lire → ajouter → vérifier.
from __future__ import annotations
from datetime import datetime
import os
import daft
import lance_namespace
import pyarrow as pa
from lance_namespace import CreateTableRequest, DescribeTableRequest
import lance_dlf # noqa: F401
CONFIG = {
"uri": "http://<DLF-ENDPOINT>",
"warehouse": "<YOUR-CATALOG>",
"token.provider": "dlf",
"dlf.region": "<REGION-ID>",
"dlf.access-key-id": "<ACCESS-KEY-ID>",
"dlf.access-key-secret": "<ACCESS-KEY-SECRET>",
"dlf.oss-endpoint": "<OSS-ENDPOINT>", # Required only for public network access to DLF
}
DATABASE = "default"
# Serialize an Arrow table to IPC bytes
def arrow_table_to_ipc_bytes(table: pa.Table) -> bytes:
sink = pa.BufferOutputStream()
with pa.ipc.new_stream(sink, table.schema) as writer:
writer.write_table(table)
return sink.getvalue().to_pybytes()
# Set OSS credential environment variables
def apply_oss_environment(storage_options: dict) -> None:
os.environ["OSS_ENDPOINT"] = storage_options["oss_endpoint"]
os.environ["OSS_ACCESS_KEY_ID"] = storage_options["oss_access_key_id"]
os.environ["OSS_ACCESS_KEY_SECRET"] = storage_options["oss_secret_access_key"]
if storage_options.get("oss_security_token"):
os.environ["OSS_SECURITY_TOKEN"] = storage_options["oss_security_token"]
if storage_options.get("oss_region"):
os.environ["OSS_REGION"] = storage_options["oss_region"]
def df_to_pydict(df):
try:
return df.to_pydict()
except AttributeError:
return df.collect().to_pydict()
def main() -> None:
ns = lance_namespace.connect("dlf", CONFIG)
# 1. Create a new table
table_name = "test_lance_daft_" + datetime.now().strftime("%Y%m%d_%H%M%S")
table_id = [DATABASE, table_name]
rows = {
"f0": [201, 202, 203],
"f1": ["daft-a", "daft-b", "daft-c"],
}
arrow_table = pa.table(rows)
create_response = ns.create_table(
CreateTableRequest(id=table_id),
arrow_table_to_ipc_bytes(arrow_table),
)
print("created:", ".".join(table_id))
print("location:", create_response.location)
# 2. Get credentials
desc = ns.describe_table(DescribeTableRequest(id=table_id))
apply_oss_environment(desc.storage_options or {})
# 3. Read and verify
read_df = daft.read_lance(desc.location)
read_df.show()
if df_to_pydict(read_df) != rows:
raise AssertionError("Initial readback mismatch")
# 4. Append data
append_rows = {
"f0": [204],
"f1": ["daft-d"],
}
append_df = daft.from_pydict(append_rows)
append_df.write_lance(desc.location, mode="append").show()
# 5. Final verification
appended_df = daft.read_lance(desc.location)
appended_df.show()
expected = {
"f0": rows["f0"] + append_rows["f0"],
"f1": rows["f1"] + append_rows["f1"],
}
if df_to_pydict(appended_df) != expected:
raise AssertionError("Daft append readback mismatch")
print("daft + dlf + lance: ok")
if __name__ == "__main__":
main()
Remarques importantes
Initialisation d'une nouvelle table : Utilisez
ns.create_table(...)pour créer une nouvelle table et écrire le premier lot de données. Utilisez Daft pour toutes les lectures et écritures suivantes.Sanitisation des journaux : Le dictionnaire
storage_optionscomplet contient des valeurs temporaires AK/SK/token. Pour des raisons de sécurité, affichez uniquement la liste des clés :print(sorted((desc.storage_options or {}).keys()))