yield per document during rebuild for live progress reporting
This commit is contained in:
parent
a45d82f206
commit
ffc7b95375
1 changed files with 8 additions and 24 deletions
|
|
@ -250,7 +250,6 @@ async def _rebuild_rechunk(
|
||||||
|
|
||||||
pending_chunks: list[Chunk] = []
|
pending_chunks: list[Chunk] = []
|
||||||
pending_docs: list[Document] = []
|
pending_docs: list[Document] = []
|
||||||
pending_doc_ids: list[str] = []
|
|
||||||
embedder = get_embedder(client._config)
|
embedder = get_embedder(client._config)
|
||||||
|
|
||||||
async for doc in _hydrate(client, documents):
|
async for doc in _hydrate(client, documents):
|
||||||
|
|
@ -283,22 +282,22 @@ async def _rebuild_rechunk(
|
||||||
|
|
||||||
pending_chunks.extend(embedded_chunks)
|
pending_chunks.extend(embedded_chunks)
|
||||||
pending_docs.append(doc)
|
pending_docs.append(doc)
|
||||||
pending_doc_ids.append(doc.id)
|
# Yield per-doc so progress reporting moves immediately. The actual
|
||||||
|
# write batches up to _REBUILD_BATCH_SIZE for throughput; if the
|
||||||
|
# process is interrupted between yield and flush, up to one
|
||||||
|
# batch's worth of trailing yields aren't persisted, which is
|
||||||
|
# consistent with the rebuild already being non-atomic.
|
||||||
|
yield doc.id
|
||||||
|
|
||||||
# Flush batch when size reached
|
# Flush batch when size reached
|
||||||
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
pending_chunks = []
|
pending_chunks = []
|
||||||
pending_docs = []
|
pending_docs = []
|
||||||
pending_doc_ids = []
|
|
||||||
|
|
||||||
# Flush remaining
|
# Flush remaining
|
||||||
if pending_docs:
|
if pending_docs:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
|
|
||||||
|
|
||||||
async def _patch_picture_descriptions(client: "HaikuRAG", doc: Document) -> int:
|
async def _patch_picture_descriptions(client: "HaikuRAG", doc: Document) -> int:
|
||||||
|
|
@ -384,7 +383,6 @@ async def _rebuild_descriptions(
|
||||||
|
|
||||||
pending_chunks: list[Chunk] = []
|
pending_chunks: list[Chunk] = []
|
||||||
pending_docs: list[Document] = []
|
pending_docs: list[Document] = []
|
||||||
pending_doc_ids: list[str] = []
|
|
||||||
embedder = get_embedder(client._config)
|
embedder = get_embedder(client._config)
|
||||||
|
|
||||||
described_total = 0
|
described_total = 0
|
||||||
|
|
@ -421,20 +419,15 @@ async def _rebuild_descriptions(
|
||||||
|
|
||||||
pending_chunks.extend(embedded_chunks)
|
pending_chunks.extend(embedded_chunks)
|
||||||
pending_docs.append(doc)
|
pending_docs.append(doc)
|
||||||
pending_doc_ids.append(doc.id)
|
yield doc.id
|
||||||
|
|
||||||
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
pending_chunks = []
|
pending_chunks = []
|
||||||
pending_docs = []
|
pending_docs = []
|
||||||
pending_doc_ids = []
|
|
||||||
|
|
||||||
if pending_docs:
|
if pending_docs:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"rebuild --descriptions: %d new picture descriptions added across %d documents",
|
"rebuild --descriptions: %d new picture descriptions added across %d documents",
|
||||||
|
|
@ -451,7 +444,6 @@ async def _rebuild_full(
|
||||||
|
|
||||||
pending_chunks: list[Chunk] = []
|
pending_chunks: list[Chunk] = []
|
||||||
pending_docs: list[Document] = []
|
pending_docs: list[Document] = []
|
||||||
pending_doc_ids: list[str] = []
|
|
||||||
converter = get_converter(client._config)
|
converter = get_converter(client._config)
|
||||||
|
|
||||||
for light_doc in documents:
|
for light_doc in documents:
|
||||||
|
|
@ -464,11 +456,8 @@ async def _rebuild_full(
|
||||||
# Flush pending batch before source rebuild (creates new doc)
|
# Flush pending batch before source rebuild (creates new doc)
|
||||||
if pending_docs:
|
if pending_docs:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
pending_chunks = []
|
pending_chunks = []
|
||||||
pending_docs = []
|
pending_docs = []
|
||||||
pending_doc_ids = []
|
|
||||||
|
|
||||||
await client.delete_document(light_doc.id)
|
await client.delete_document(light_doc.id)
|
||||||
new_doc = await client.create_document_from_source(
|
new_doc = await client.create_document_from_source(
|
||||||
|
|
@ -508,19 +497,14 @@ async def _rebuild_full(
|
||||||
|
|
||||||
pending_chunks.extend(embedded_chunks)
|
pending_chunks.extend(embedded_chunks)
|
||||||
pending_docs.append(doc)
|
pending_docs.append(doc)
|
||||||
pending_doc_ids.append(doc.id)
|
yield doc.id
|
||||||
|
|
||||||
# Flush batch when size reached
|
# Flush batch when size reached
|
||||||
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
if len(pending_docs) >= _REBUILD_BATCH_SIZE:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
pending_chunks = []
|
pending_chunks = []
|
||||||
pending_docs = []
|
pending_docs = []
|
||||||
pending_doc_ids = []
|
|
||||||
|
|
||||||
# Flush remaining
|
# Flush remaining
|
||||||
if pending_docs:
|
if pending_docs:
|
||||||
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
await _flush_rebuild_batch(client, pending_docs, pending_chunks)
|
||||||
for doc_id in pending_doc_ids:
|
|
||||||
yield doc_id
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue