Skip to content
Merged
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
9 changes: 0 additions & 9 deletions java/lance-jni/src/mem_wal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1265,9 +1265,6 @@ fn build_writer_config(env: &mut JNIEnv, config: &JObject) -> Result<ShardWriter
if let Some(v) = read_optional_bool(env, config, "durableWrite")? {
writer_config = writer_config.with_durable_write(v);
}
if let Some(v) = read_optional_bool(env, config, "syncIndexedWrite")? {
writer_config = writer_config.with_sync_indexed_write(v);
}
if let Some(v) = read_optional_u64(env, config, "maxWalBufferSize")? {
writer_config = writer_config.with_max_wal_buffer_size(v as usize);
}
Expand All @@ -1289,12 +1286,6 @@ fn build_writer_config(env: &mut JNIEnv, config: &JObject) -> Result<ShardWriter
if let Some(v) = read_optional_u64(env, config, "manifestScanBatchSize")? {
writer_config = writer_config.with_manifest_scan_batch_size(v as usize);
}
if let Some(v) = read_optional_u64(env, config, "asyncIndexBufferRows")? {
writer_config = writer_config.with_async_index_buffer_rows(v as usize);
}
if let Some(v) = read_optional_u64(env, config, "asyncIndexIntervalMs")? {
writer_config = writer_config.with_async_index_interval(Duration::from_millis(v));
}
if let Some(v) = read_optional_u64(env, config, "backpressureLogIntervalMs")? {
writer_config = writer_config.with_backpressure_log_interval(Duration::from_millis(v));
}
Expand Down
41 changes: 0 additions & 41 deletions java/src/main/java/org/lance/memwal/ShardWriterConfig.java
Original file line number Diff line number Diff line change
Expand Up @@ -28,16 +28,13 @@
*/
public class ShardWriterConfig {
private Optional<Boolean> durableWrite = Optional.empty();
private Optional<Boolean> syncIndexedWrite = Optional.empty();
private Optional<Long> maxWalBufferSize = Optional.empty();
private Optional<Long> maxWalFlushIntervalMs = Optional.empty();
private Optional<Long> maxMemtableSize = Optional.empty();
private Optional<Long> maxMemtableRows = Optional.empty();
private Optional<Long> maxMemtableBatches = Optional.empty();
private Optional<Long> maxUnflushedMemtableBytes = Optional.empty();
private Optional<Long> manifestScanBatchSize = Optional.empty();
private Optional<Long> asyncIndexBufferRows = Optional.empty();
private Optional<Long> asyncIndexIntervalMs = Optional.empty();
private Optional<Long> backpressureLogIntervalMs = Optional.empty();
private Optional<Long> statsLogIntervalMs = Optional.empty();
private List<MemWalHnswParams> hnswParams = Collections.emptyList();
Expand All @@ -48,12 +45,6 @@ public ShardWriterConfig withDurableWrite(boolean durableWrite) {
return this;
}

/** Whether indexed writes are applied synchronously. */
public ShardWriterConfig withSyncIndexedWrite(boolean syncIndexedWrite) {
this.syncIndexedWrite = Optional.of(syncIndexedWrite);
return this;
}

/** Maximum size of the in-memory WAL buffer, in bytes. */
public ShardWriterConfig withMaxWalBufferSize(long maxWalBufferSize) {
Preconditions.checkArgument(
Expand Down Expand Up @@ -118,26 +109,6 @@ public ShardWriterConfig withManifestScanBatchSize(long manifestScanBatchSize) {
return this;
}

/** Number of rows buffered before an asynchronous index update is triggered. */
public ShardWriterConfig withAsyncIndexBufferRows(long asyncIndexBufferRows) {
Preconditions.checkArgument(
asyncIndexBufferRows >= 0,
"asyncIndexBufferRows must not be negative, got %s",
asyncIndexBufferRows);
this.asyncIndexBufferRows = Optional.of(asyncIndexBufferRows);
return this;
}

/** Interval between asynchronous index updates, in milliseconds. */
public ShardWriterConfig withAsyncIndexIntervalMs(long asyncIndexIntervalMs) {
Preconditions.checkArgument(
asyncIndexIntervalMs >= 0,
"asyncIndexIntervalMs must not be negative, got %s",
asyncIndexIntervalMs);
this.asyncIndexIntervalMs = Optional.of(asyncIndexIntervalMs);
return this;
}

/** Interval between backpressure log messages, in milliseconds. */
public ShardWriterConfig withBackpressureLogIntervalMs(long backpressureLogIntervalMs) {
Preconditions.checkArgument(
Expand Down Expand Up @@ -172,10 +143,6 @@ public Optional<Boolean> durableWrite() {
return durableWrite;
}

public Optional<Boolean> syncIndexedWrite() {
return syncIndexedWrite;
}

public Optional<Long> maxWalBufferSize() {
return maxWalBufferSize;
}
Expand Down Expand Up @@ -204,14 +171,6 @@ public Optional<Long> manifestScanBatchSize() {
return manifestScanBatchSize;
}

public Optional<Long> asyncIndexBufferRows() {
return asyncIndexBufferRows;
}

public Optional<Long> asyncIndexIntervalMs() {
return asyncIndexIntervalMs;
}

public Optional<Long> backpressureLogIntervalMs() {
return backpressureLogIntervalMs;
}
Expand Down
1 change: 0 additions & 1 deletion java/src/test/java/org/lance/memwal/MemWalTest.java
Original file line number Diff line number Diff line change
Expand Up @@ -409,7 +409,6 @@ void testShardWriterDeleteMasksBaseRow(@TempDir Path tempDir) throws Exception {
ShardWriterConfig config =
new ShardWriterConfig()
.withDurableWrite(true)
.withSyncIndexedWrite(true)
.withMaxWalBufferSize(1)
.withMaxWalFlushIntervalMs(10);

Expand Down
18 changes: 0 additions & 18 deletions python/python/lance/dataset.py
Original file line number Diff line number Diff line change
Expand Up @@ -5150,16 +5150,13 @@ def initialize_mem_wal(
identity_column: Optional[str] = None,
unsharded: bool = False,
durable_write: Optional[bool] = None,
sync_indexed_write: Optional[bool] = None,
max_wal_buffer_size: Optional[int] = None,
max_wal_flush_interval_ms: Optional[int] = None,
max_memtable_size: Optional[int] = None,
max_memtable_rows: Optional[int] = None,
max_memtable_batches: Optional[int] = None,
max_unflushed_memtable_bytes: Optional[int] = None,
manifest_scan_batch_size: Optional[int] = None,
async_index_buffer_rows: Optional[int] = None,
async_index_interval_ms: Optional[int] = None,
backpressure_log_interval_ms: Optional[int] = None,
stats_log_interval_ms: Optional[int] = None,
hnsw_params: Optional[Dict[str, Dict[str, int]]] = None,
Expand Down Expand Up @@ -5220,16 +5217,13 @@ def initialize_mem_wal(
identity_column=identity_column,
unsharded=unsharded,
durable_write=durable_write,
sync_indexed_write=sync_indexed_write,
max_wal_buffer_size=max_wal_buffer_size,
max_wal_flush_interval_ms=max_wal_flush_interval_ms,
max_memtable_size=max_memtable_size,
max_memtable_rows=max_memtable_rows,
max_memtable_batches=max_memtable_batches,
max_unflushed_memtable_bytes=max_unflushed_memtable_bytes,
manifest_scan_batch_size=manifest_scan_batch_size,
async_index_buffer_rows=async_index_buffer_rows,
async_index_interval_ms=async_index_interval_ms,
backpressure_log_interval_ms=backpressure_log_interval_ms,
stats_log_interval_ms=stats_log_interval_ms,
hnsw_params=hnsw_params,
Expand All @@ -5252,16 +5246,13 @@ def mem_wal_writer(
shard_id: str,
*,
durable_write: Optional[bool] = None,
sync_indexed_write: Optional[bool] = None,
max_wal_buffer_size: Optional[int] = None,
max_wal_flush_interval_ms: Optional[int] = None,
max_memtable_size: Optional[int] = None,
max_memtable_rows: Optional[int] = None,
max_memtable_batches: Optional[int] = None,
max_unflushed_memtable_bytes: Optional[int] = None,
manifest_scan_batch_size: Optional[int] = None,
async_index_buffer_rows: Optional[int] = None,
async_index_interval_ms: Optional[int] = None,
backpressure_log_interval_ms: Optional[int] = None,
stats_log_interval_ms: Optional[int] = None,
hnsw_params: Optional[Dict[str, Dict[str, int]]] = None,
Expand All @@ -5279,8 +5270,6 @@ def mem_wal_writer(
``str(uuid.uuid4())``).
durable_write : bool, optional
Whether to fsync WAL writes (default: ``True``).
sync_indexed_write : bool, optional
Whether index updates are synchronous (default: ``True``).
max_wal_buffer_size : int, optional
Maximum WAL buffer size in bytes (default: 10 MB).
max_wal_flush_interval_ms : int, optional
Expand All @@ -5295,10 +5284,6 @@ def mem_wal_writer(
Maximum unflushed bytes before backpressure (default: 1 GB).
manifest_scan_batch_size : int, optional
Batch size for manifest scans (default: 2).
async_index_buffer_rows : int, optional
Buffer rows for async index updates (default: 10 000).
async_index_interval_ms : int, optional
Interval for async index updates in milliseconds (default: 1000).
backpressure_log_interval_ms : int, optional
Interval for backpressure log messages in milliseconds
(default: 30 000).
Expand Down Expand Up @@ -5346,16 +5331,13 @@ def mem_wal_writer(
name: val
for name, val in [
("durable_write", durable_write),
("sync_indexed_write", sync_indexed_write),
("max_wal_buffer_size", max_wal_buffer_size),
("max_wal_flush_interval_ms", max_wal_flush_interval_ms),
("max_memtable_size", max_memtable_size),
("max_memtable_rows", max_memtable_rows),
("max_memtable_batches", max_memtable_batches),
("max_unflushed_memtable_bytes", max_unflushed_memtable_bytes),
("manifest_scan_batch_size", manifest_scan_batch_size),
("async_index_buffer_rows", async_index_buffer_rows),
("async_index_interval_ms", async_index_interval_ms),
("backpressure_log_interval_ms", backpressure_log_interval_ms),
("stats_log_interval_ms", stats_log_interval_ms),
("hnsw_params", hnsw_params),
Expand Down
7 changes: 1 addition & 6 deletions python/python/tests/test_mem_wal.py
Original file line number Diff line number Diff line change
Expand Up @@ -210,7 +210,6 @@ def test_shard_writer_delete_binding_masks_base_row(tmp_path):
with ds.mem_wal_writer(
shard_id,
durable_write=True,
sync_indexed_write=True,
max_wal_buffer_size=1,
max_wal_flush_interval_ms=10,
) as writer:
Expand Down Expand Up @@ -356,7 +355,6 @@ def test_shard_writer_e2e_correctness(tmp_path):
writer = ds.mem_wal_writer(
shard_id,
durable_write=True,
sync_indexed_write=True,
max_wal_buffer_size=10 * 1024, # 10 KB
max_wal_flush_interval_ms=50,
max_memtable_size=80, # flush after ~80 rows
Expand Down Expand Up @@ -406,9 +404,7 @@ def test_shard_writer_e2e_correctness(tmp_path):
# === New writer: write and read back via active MemTable scanner ===
ds2 = lance.dataset(ds_path)
shard_id2 = str(uuid.uuid4())
with ds2.mem_wal_writer(
shard_id2, durable_write=False, sync_indexed_write=True
) as writer2:
with ds2.mem_wal_writer(shard_id2, durable_write=False) as writer2:
verify_batch = _e2e_batch(schema, start_id=10000, num_rows=10)
writer2.put(pa.Table.from_batches([verify_batch]))
result = writer2.lsm_scanner().to_table()
Expand Down Expand Up @@ -523,7 +519,6 @@ def test_initialize_mem_wal_writer_config_defaults(tmp_path):
# Duration knobs are recorded in milliseconds with a `_ms` suffix.
assert defaults["max_wal_flush_interval_ms"] == "250"
# Every ShardWriterConfig tunable is recorded once any default is set.
assert "sync_indexed_write" in defaults
assert "enable_memtable" in defaults


Expand Down
33 changes: 0 additions & 33 deletions python/src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -203,16 +203,13 @@ fn stats_log_interval_from_millis(ms: u64) -> Option<std::time::Duration> {
#[allow(clippy::too_many_arguments)]
fn writer_config_from_kwargs(
durable_write: Option<bool>,
sync_indexed_write: Option<bool>,
max_wal_buffer_size: Option<usize>,
max_wal_flush_interval_ms: Option<u64>,
max_memtable_size: Option<usize>,
max_memtable_rows: Option<usize>,
max_memtable_batches: Option<usize>,
max_unflushed_memtable_bytes: Option<usize>,
manifest_scan_batch_size: Option<usize>,
async_index_buffer_rows: Option<usize>,
async_index_interval_ms: Option<u64>,
backpressure_log_interval_ms: Option<u64>,
stats_log_interval_ms: Option<u64>,
hnsw_params: Option<HashMap<String, HashMap<String, u32>>>,
Expand All @@ -225,10 +222,6 @@ fn writer_config_from_kwargs(
config = config.with_durable_write(v);
any = true;
}
if let Some(v) = sync_indexed_write {
config = config.with_sync_indexed_write(v);
any = true;
}
if let Some(v) = max_wal_buffer_size {
config = config.with_max_wal_buffer_size(v);
any = true;
Expand Down Expand Up @@ -257,14 +250,6 @@ fn writer_config_from_kwargs(
config = config.with_manifest_scan_batch_size(v);
any = true;
}
if let Some(v) = async_index_buffer_rows {
config = config.with_async_index_buffer_rows(v);
any = true;
}
if let Some(v) = async_index_interval_ms {
config = config.with_async_index_interval(Duration::from_millis(v));
any = true;
}
if let Some(v) = backpressure_log_interval_ms {
config = config.with_backpressure_log_interval(Duration::from_millis(v));
any = true;
Expand Down Expand Up @@ -3586,16 +3571,13 @@ impl Dataset {
identity_column=None,
unsharded=false,
durable_write=None,
sync_indexed_write=None,
max_wal_buffer_size=None,
max_wal_flush_interval_ms=None,
max_memtable_size=None,
max_memtable_rows=None,
max_memtable_batches=None,
max_unflushed_memtable_bytes=None,
manifest_scan_batch_size=None,
async_index_buffer_rows=None,
async_index_interval_ms=None,
backpressure_log_interval_ms=None,
stats_log_interval_ms=None,
hnsw_params=None,
Expand All @@ -3609,16 +3591,13 @@ impl Dataset {
identity_column: Option<String>,
unsharded: bool,
durable_write: Option<bool>,
sync_indexed_write: Option<bool>,
max_wal_buffer_size: Option<usize>,
max_wal_flush_interval_ms: Option<u64>,
max_memtable_size: Option<usize>,
max_memtable_rows: Option<usize>,
max_memtable_batches: Option<usize>,
max_unflushed_memtable_bytes: Option<usize>,
manifest_scan_batch_size: Option<usize>,
async_index_buffer_rows: Option<usize>,
async_index_interval_ms: Option<u64>,
backpressure_log_interval_ms: Option<u64>,
stats_log_interval_ms: Option<u64>,
hnsw_params: Option<HashMap<String, HashMap<String, u32>>>,
Expand Down Expand Up @@ -3646,16 +3625,13 @@ impl Dataset {

let writer_config = writer_config_from_kwargs(
durable_write,
sync_indexed_write,
max_wal_buffer_size,
max_wal_flush_interval_ms,
max_memtable_size,
max_memtable_rows,
max_memtable_batches,
max_unflushed_memtable_bytes,
manifest_scan_batch_size,
async_index_buffer_rows,
async_index_interval_ms,
backpressure_log_interval_ms,
stats_log_interval_ms,
hnsw_params,
Expand Down Expand Up @@ -3737,16 +3713,13 @@ impl Dataset {
shard_id,
*,
durable_write=None,
sync_indexed_write=None,
max_wal_buffer_size=None,
max_wal_flush_interval_ms=None,
max_memtable_size=None,
max_memtable_rows=None,
max_memtable_batches=None,
max_unflushed_memtable_bytes=None,
manifest_scan_batch_size=None,
async_index_buffer_rows=None,
async_index_interval_ms=None,
backpressure_log_interval_ms=None,
stats_log_interval_ms=None,
hnsw_params=None,
Expand All @@ -3756,16 +3729,13 @@ impl Dataset {
py: Python<'_>,
shard_id: String,
durable_write: Option<bool>,
sync_indexed_write: Option<bool>,
max_wal_buffer_size: Option<usize>,
max_wal_flush_interval_ms: Option<u64>,
max_memtable_size: Option<usize>,
max_memtable_rows: Option<usize>,
max_memtable_batches: Option<usize>,
max_unflushed_memtable_bytes: Option<usize>,
manifest_scan_batch_size: Option<usize>,
async_index_buffer_rows: Option<usize>,
async_index_interval_ms: Option<u64>,
backpressure_log_interval_ms: Option<u64>,
stats_log_interval_ms: Option<u64>,
hnsw_params: Option<HashMap<String, HashMap<String, u32>>>,
Expand All @@ -3777,16 +3747,13 @@ impl Dataset {

let config = writer_config_from_kwargs(
durable_write,
sync_indexed_write,
max_wal_buffer_size,
max_wal_flush_interval_ms,
max_memtable_size,
max_memtable_rows,
max_memtable_batches,
max_unflushed_memtable_bytes,
manifest_scan_batch_size,
async_index_buffer_rows,
async_index_interval_ms,
backpressure_log_interval_ms,
stats_log_interval_ms,
hnsw_params,
Expand Down
Loading
Loading