Compare commits

..

1 Commits

Author SHA1 Message Date
lancedb automation 4419664740 chore: update lance dependency to v12.0.0-beta.10 2026-09-01 16:19:40 +00:00
14 changed files with 70 additions and 310 deletions
Generated
+42 -42
View File
@@ -3455,8 +3455,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
[[package]]
name = "fsst"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"rand 0.9.5",
@@ -4815,8 +4815,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
[[package]]
name = "lance"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arc-swap",
"arrow",
@@ -4888,8 +4888,8 @@ dependencies = [
[[package]]
name = "lance-arrow"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4911,7 +4911,7 @@ dependencies = [
[[package]]
name = "lance-arrow-scalar"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4925,7 +4925,7 @@ dependencies = [
[[package]]
name = "lance-arrow-stats"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -4934,8 +4934,8 @@ dependencies = [
[[package]]
name = "lance-bitpacking"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrayref",
"crunchy",
@@ -4945,8 +4945,8 @@ dependencies = [
[[package]]
name = "lance-core"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4983,8 +4983,8 @@ dependencies = [
[[package]]
name = "lance-datafusion"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"arrow-array",
@@ -5014,8 +5014,8 @@ dependencies = [
[[package]]
name = "lance-datagen"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"arrow-array",
@@ -5032,8 +5032,8 @@ dependencies = [
[[package]]
name = "lance-derive"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"proc-macro2",
"quote",
@@ -5042,8 +5042,8 @@ dependencies = [
[[package]]
name = "lance-encoding"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5076,8 +5076,8 @@ dependencies = [
[[package]]
name = "lance-file"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5108,8 +5108,8 @@ dependencies = [
[[package]]
name = "lance-index"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arc-swap",
"arrow",
@@ -5173,8 +5173,8 @@ dependencies = [
[[package]]
name = "lance-index-core"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5196,8 +5196,8 @@ dependencies = [
[[package]]
name = "lance-io"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"arrow-array",
@@ -5237,8 +5237,8 @@ dependencies = [
[[package]]
name = "lance-linalg"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5252,8 +5252,8 @@ dependencies = [
[[package]]
name = "lance-namespace"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"async-trait",
@@ -5265,8 +5265,8 @@ dependencies = [
[[package]]
name = "lance-namespace-impls"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"arrow-ipc",
@@ -5319,8 +5319,8 @@ dependencies = [
[[package]]
name = "lance-select"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5334,8 +5334,8 @@ dependencies = [
[[package]]
name = "lance-table"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow",
"arrow-array",
@@ -5375,8 +5375,8 @@ dependencies = [
[[package]]
name = "lance-testing"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5389,8 +5389,8 @@ dependencies = [
[[package]]
name = "lance-tokenizer"
version = "12.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.9#6f93e3fe389f5a661a02aff00c14df6be7bf0505"
version = "12.0.0-beta.10"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.10#93550f3eec61c2f1b23f605fa38914d7ac5a526e"
dependencies = [
"frostem",
"icu_segmenter",
+14 -14
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=12.0.0-beta.9", default-features = false, "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=12.0.0-beta.9", default-features = false, "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=12.0.0-beta.9", default-features = false, "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=12.0.0-beta.9", "tag" = "v12.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=12.0.0-beta.10", default-features = false, "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=12.0.0-beta.10", default-features = false, "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=12.0.0-beta.10", default-features = false, "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=12.0.0-beta.10", "tag" = "v12.0.0-beta.10", "git" = "https://github.com/lance-format/lance.git" }
lancedb = { path = "rust/lancedb", default-features = false }
ahash = "0.8"
# Note that this one does not include pyarrow
-9
View File
@@ -446,15 +446,6 @@ paths:
properties:
column:
type: string
name:
type: string
description: Optional name for the created index.
replace:
type: boolean
default: true
description: |
Whether to replace an existing index with the same resolved
name. Defaults to true.
metric_type:
type: string
nullable: false
-8
View File
@@ -676,17 +676,9 @@ List all the versions of the table
abstract mergeInsert(on): MergeInsertBuilder
```
Create a [MergeInsertBuilder](MergeInsertBuilder.md), which combines new data with the
existing table in a single transaction — inserting, updating and deleting
rows depending on how they match.
#### Parameters
* **on**: `string` \| `string`[]
The column, or columns, to match source rows against target
rows on. Typically a key or id column. Several columns match on the
composite key: a source row updates a target row only when it agrees on
every one of them.
#### Returns
+1 -1
View File
@@ -28,7 +28,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version>
<lance-core.version>12.0.0-beta.9</lance-core.version>
<lance-core.version>12.0.0-beta.10</lance-core.version>
<spotless.skip>false</spotless.skip>
<spotless.version>2.30.0</spotless.version>
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
+1 -34
View File
@@ -737,12 +737,11 @@ it("should query documents with LangChain PDF metadata", async () => {
describe("merge insert", () => {
let tmpDir: tmp.DirResult;
let conn: Connection;
let table: Table;
beforeEach(async () => {
tmpDir = tmp.dirSync({ unsafeCleanup: true });
conn = await connect(tmpDir.name);
const conn = await connect(tmpDir.name);
table = await conn.createTable("some_table", [
{ a: 1, b: "a" },
@@ -780,38 +779,6 @@ describe("merge insert", () => {
expect(result.map((row) => ({ ...row }))).toEqual(expected);
});
test("upsert on a composite key", async () => {
const composite = await conn.createTable("composite", [
{ shard: "a", id: 1, val: "x" },
{ shard: "a", id: 2, val: "y" },
{ shard: "b", id: 1, val: "z" },
]);
// ("a", 1) matches an existing row and updates it. ("b", 2) agrees with an
// existing row on each key column separately but on neither pair, so it is
// an insert.
const mergeInsertRes = await composite
.mergeInsert(["shard", "id"])
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute([
{ shard: "a", id: 1, val: "X" },
{ shard: "b", id: 2, val: "W" },
]);
expect(mergeInsertRes.numUpdatedRows).toBe(1);
expect(mergeInsertRes.numInsertedRows).toBe(1);
const result = (await composite.toArrow())
.toArray()
.sort((a, b) => a.shard.localeCompare(b.shard) || a.id - b.id);
expect(result.map((row) => ({ ...row }))).toEqual([
{ shard: "a", id: 1, val: "X" },
{ shard: "a", id: 2, val: "y" },
{ shard: "b", id: 1, val: "z" },
{ shard: "b", id: 2, val: "W" },
]);
});
test("conditional update", async () => {
const newData = [
{ a: 2, b: "x" },
-10
View File
@@ -919,16 +919,6 @@ export abstract class Table {
/** Return the table as an arrow table */
abstract toArrow(): Promise<ArrowTable>;
/**
* Create a {@link MergeInsertBuilder}, which combines new data with the
* existing table in a single transaction — inserting, updating and deleting
* rows depending on how they match.
*
* @param on - The column, or columns, to match source rows against target
* rows on. Typically a key or id column. Several columns match on the
* composite key: a source row updates a target row only when it agrees on
* every one of them.
*/
abstract mergeInsert(on: string | string[]): MergeInsertBuilder;
/** List all the stats of a specified index
-1
View File
@@ -548,7 +548,6 @@ class RemoteTable(Table):
LOOP.run(
self._table.create_index(
column,
replace=replace,
config=config,
wait_timeout=wait_timeout,
name=name,
+2 -6
View File
@@ -1547,9 +1547,7 @@ class Table(ABC):
on: Union[str, Iterable[str]]
A column (or columns) to join on. This is how records from the
source table and target table are matched. Typically this is some
kind of key or id column. Passing several columns matches on the
composite key: a source row updates a target row only when it
agrees on every one of them.
kind of key or id column.
Examples
--------
@@ -5703,9 +5701,7 @@ class AsyncTable:
on: Union[str, Iterable[str]]
A column (or columns) to join on. This is how records from the
source table and target table are matched. Typically this is some
kind of key or id column. Passing several columns matches on the
composite key: a source row updates a target row only when it
agrees on every one of them.
kind of key or id column.
Examples
--------
+1 -9
View File
@@ -820,13 +820,11 @@ def test_table_create_indices():
scalar_req = received_requests[0]
assert "name" in scalar_req
assert scalar_req["name"] == "custom_scalar_idx"
assert scalar_req["replace"] is False
# Check FTS index request has custom name
fts_req = received_requests[1]
assert "name" in fts_req
assert fts_req["name"] == "custom_fts_idx"
assert fts_req["replace"] is False
assert fts_req["block_size"] == 256
assert fts_req["custom_stop_words"] == ["cloud"]
@@ -834,7 +832,6 @@ def test_table_create_indices():
vector_req = received_requests[2]
assert "name" in vector_req
assert vector_req["name"] == "custom_vector_idx"
assert "replace" not in vector_req
table.wait_for_index(["custom_scalar_idx"], timedelta(seconds=2))
table.wait_for_index(
@@ -1107,9 +1104,6 @@ def test_remote_create_index_new_api():
table.create_index("text", config=FTS(block_size=256))
# IvfRq via new API
table.create_index("vector", config=IvfRq(distance_type="l2"))
table.create_index(
"vector", config=IvfPq(distance_type="l2"), replace=False
)
# Legacy index_type="IVF_RQ" routes to IvfRq config under the hood.
with pytest.warns(DeprecationWarning, match="create_index"):
@@ -1119,17 +1113,15 @@ def test_remote_create_index_new_api():
num_partitions=8,
)
assert len(received_requests) == 6
assert len(received_requests) == 5
assert [req["column"] for req in received_requests] == [
"vector",
"category",
"text",
"vector",
"vector",
"vector",
]
assert received_requests[2]["block_size"] == 256
assert received_requests[4]["replace"] is False
def test_table_wait_for_index_timeout():
-37
View File
@@ -2682,43 +2682,6 @@ def test_merge_insert(mem_db: DBConnection):
)
def test_merge_insert_composite_key(mem_db: DBConnection):
table = mem_db.create_table(
"my_table",
data=pa.table(
{
"shard": ["a", "a", "b"],
"id": [1, 2, 1],
"val": ["x", "y", "z"],
}
),
)
# ("a", 1) matches an existing row and updates it. ("b", 2) agrees with an
# existing row on each key column separately but on neither pair, so it is
# an insert.
new_data = pa.table({"shard": ["a", "b"], "id": [1, 2], "val": ["X", "W"]})
res = (
table.merge_insert(["shard", "id"])
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(new_data)
)
assert res.num_updated_rows == 1
assert res.num_inserted_rows == 1
expected = pa.table(
{
"shard": ["a", "a", "b", "b"],
"id": [1, 2, 1, 2],
"val": ["X", "y", "z", "W"],
}
)
assert table.to_arrow().sort_by([("shard", "ascending"), ("id", "ascending")]) == (
expected
)
def test_merge_insert_nullable_pandas_into_pydantic_schema(mem_db: DBConnection):
# Regression test for https://github.com/lancedb/lancedb/issues/2366
pd = pytest.importorskip("pandas")
+7 -134
View File
@@ -72,7 +72,7 @@ use lance_datafusion::exec::{OneShotExec, execute_plan};
use reqwest::{RequestBuilder, Response};
use serde::{Deserialize, Serialize};
use serde_json::Number;
use std::collections::{HashMap, HashSet};
use std::collections::HashMap;
use std::io::Cursor;
use std::pin::Pin;
use std::sync::{Arc, Mutex};
@@ -527,10 +527,6 @@ impl<S: HttpSend> RemoteTable<S> {
"column": canonical_column
});
if !index.replace {
body["replace"] = false.into();
}
// Add name parameter if provided (for backwards compatibility, only include if Some)
if let Some(ref name) = index.name {
body["name"] = serde_json::Value::String(name.clone());
@@ -3651,12 +3647,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
#[derive(Serialize, Clone, Debug)]
pub struct MergeInsertRequest {
// Sent as one repeated `on` query parameter per column, which is how the
// namespace spec encodes an array-valued `on`. serde_urlencoded (which
// reqwest's `query()` uses) cannot serialize a sequence nested in a struct,
// so this field is emitted separately by [`Self::on_query_params`].
#[serde(skip_serializing)]
on: Vec<String>,
on: String,
when_matched_update_all: bool,
when_matched_update_all_filt: Option<String>,
when_not_matched_insert_all: bool,
@@ -3672,17 +3663,6 @@ pub struct MergeInsertRequest {
use_lsm: Option<bool>,
}
impl MergeInsertRequest {
/// The `on` columns as repeated query parameters: `?on=a&on=b`.
///
/// A single column serializes to `?on=a`, exactly what clients sent before
/// `on` became a list, so a server that predates composite keys sees no
/// change from a single-column caller.
pub(crate) fn on_query_params(&self) -> Vec<(&str, &str)> {
self.on.iter().map(|col| ("on", col.as_str())).collect()
}
}
fn is_true(b: &bool) -> bool {
*b
}
@@ -3695,15 +3675,12 @@ impl TryFrom<MergeInsertBuilder> for MergeInsertRequest {
return Err(Error::InvalidInput {
message: "MergeInsertBuilder missing required 'on' field".into(),
});
}
// The server rejects a repeated column with a 400; catching it here
// names the offending column and costs no round trip.
let mut seen = HashSet::with_capacity(value.on.len());
if let Some(dup) = value.on.iter().find(|col| !seen.insert(*col)) {
return Err(Error::InvalidInput {
message: format!("MergeInsertBuilder 'on' column '{dup}' is repeated"),
} else if value.on.len() > 1 {
return Err(Error::NotSupported {
message: "MergeInsertBuilder only supports a single 'on' column".into(),
});
}
let on = value.on[0].clone();
let when_matched_update_all_filt = match value.when_matched_update_all_filt {
Some(MergeFilter::Sql(sql)) => Some(sql),
@@ -3727,7 +3704,7 @@ impl TryFrom<MergeInsertBuilder> for MergeInsertRequest {
};
Ok(Self {
on: value.on,
on,
when_matched_update_all: value.when_matched_update_all,
when_matched_update_all_filt,
when_not_matched_insert_all: value.when_not_matched_insert_all,
@@ -4571,76 +4548,6 @@ mod tests {
}
}
#[tokio::test]
async fn test_merge_insert_composite_key() {
let batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
)
.unwrap();
let data: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
[Ok(batch.clone())],
batch.schema(),
));
let table = Table::new_with_handler("my_table", move |request| {
assert_eq!(request.url().path(), "/v1/table/my_table/merge_insert/");
// One repeated `on` per column, in the order the caller gave them.
let on = request
.url()
.query_pairs()
.filter(|(key, _)| key == "on")
.map(|(_, value)| value.into_owned())
.collect::<Vec<_>>();
assert_eq!(on, vec!["shard_key".to_string(), "id".to_string()]);
let params = request.url().query_pairs().collect::<HashMap<_, _>>();
assert_eq!(params["when_matched_update_all"], "true");
assert_eq!(params["when_not_matched_insert_all"], "true");
http::Response::builder()
.status(200)
.body(r#"{"version": 43, "num_deleted_rows": 0, "num_inserted_rows": 3, "num_updated_rows": 0}"#)
.unwrap()
});
let mut merge = table.merge_insert(&["shard_key", "id"]);
merge.when_matched_update_all(None);
merge.when_not_matched_insert_all();
let result = table.base_table().merge_insert(merge, data).await.unwrap();
assert_eq!(result.num_inserted_rows, 3);
}
#[tokio::test]
async fn test_merge_insert_rejects_repeated_on_column() {
let batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from(vec![1]))],
)
.unwrap();
let data: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
[Ok(batch.clone())],
batch.schema(),
));
let table = Table::new_with_handler::<&str>("my_table", |request| {
panic!("Unexpected request: {}", request.url());
});
let merge = table.merge_insert(&["id", "id"]);
let err = table
.base_table()
.merge_insert(merge, data)
.await
.unwrap_err();
assert!(
matches!(&err, Error::InvalidInput { message } if message.contains("'id' is repeated")),
"unexpected error: {err}"
);
}
#[tokio::test]
async fn test_merge_insert_retries_on_409() {
let batch = RecordBatch::try_new(
@@ -6326,40 +6233,6 @@ mod tests {
}
}
#[tokio::test]
async fn test_create_index_forwards_replace_false_on_existing_route() {
let table = Table::new_with_handler("my_table", move |request| {
assert_eq!(request.method(), "POST");
match request.url().path() {
"/v1/table/my_table/describe/" => {
let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]);
http::Response::builder()
.status(200)
.body(describe_response(&schema))
.unwrap()
}
"/v1/table/my_table/create_index/" => {
let body = request.body().unwrap().as_bytes().unwrap();
let body: serde_json::Value = serde_json::from_slice(body).unwrap();
assert_eq!(body["replace"], json!(false));
http::Response::builder()
.status(200)
.body("{}".to_string())
.unwrap()
}
path => panic!("Unexpected path: {}", path),
}
});
table
.create_index(&["a"], Index::BTree(Default::default()))
.replace(false)
.execute()
.await
.unwrap();
}
#[tokio::test]
async fn test_create_index_returns_job() {
let describe_calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
+1 -2
View File
@@ -734,7 +734,6 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteWriteExec<S> {
WriteOp::MergeInsert { query, timeout } => {
let mut request = client
.post(&format!("/v1/table/{}/merge_insert/", identifier))
.query(&query.on_query_params())
.query(query)
.header(CONTENT_TYPE, ARROW_STREAM_CONTENT_TYPE);
if let Some(timeout) = timeout {
@@ -1490,7 +1489,7 @@ mod tests {
});
let query = MergeInsertRequest {
on: vec!["id".to_string()],
on: "id".to_string(),
when_matched_update_all: false,
when_matched_update_all_filt: None,
when_not_matched_insert_all: false,
+1 -3
View File
@@ -1506,9 +1506,7 @@ impl Table {
///
/// * `on` One or more columns to join on. This is how records from the
/// source table and target table are matched. Typically this is some
/// kind of key or id column. Several columns match on the composite
/// key: a source row updates a target row only when it agrees on every
/// one of them.
/// kind of key or id column.
///
/// # Examples
///