Skip to content
163 changes: 159 additions & 4 deletions src/qlever/commands/update_wikidata.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
import re
import signal
import time
from datetime import datetime, timezone
from datetime import datetime, timedelta, timezone
from enum import Enum, auto
from pathlib import Path
from threading import Event
Expand Down Expand Up @@ -129,6 +129,9 @@ def __init__(self):
self.ctrl_c_pressed = Event()
# Set to `True` when finished.
self.finished = False
# Size of the triple history right after the last prune (see
# `prune_triple_history`).
self.triple_history_size_after_last_prune = 0

def description(self) -> str:
return "Update from given SSE stream"
Expand Down Expand Up @@ -239,6 +242,16 @@ def additional_arguments(self, subparser) -> None:
help="Process exactly this many messages and then exit "
"(default: no bound on the number of messages)",
)
subparser.add_argument(
"--triple-history-minutes",
type=int,
default=30,
help="Remember the causally newest event for each triple for "
"this many minutes of stream time, ACROSS batches, and ignore "
"inserts and deletes that are superseded by such an event; the "
"stream can reorder events by minutes, so checking only within "
"a batch is not enough (0 disables the history)",
)
subparser.add_argument(
"--verbose",
choices=["no", "yes"],
Expand Down Expand Up @@ -390,6 +403,73 @@ def determine_next_cached_update(

return cached_file_name, batch_size

def prune_triple_history(
self,
triple_history: dict[str, tuple[tuple[int, int], bool, str]],
latest_event_date: str | None,
history_minutes: int,
) -> dict[str, tuple[tuple[int, int], bool, str]]:
"""
Return the given triple history without the entries that are older
than `history_minutes` before the given latest event date (the
event dates are ISO strings and compare lexicographically).

A prune is linear in the size of the history, so it only actually
happens when the history has at least doubled in size since the
last prune; the total pruning cost is then linear in the total
number of insertions, no matter how often this is called. Keeping
entries longer than `history_minutes` is harmless for correctness,
the window is only a bound on the memory usage.
"""
if not triple_history or latest_event_date is None:
return triple_history
if len(triple_history) < max(
2 * self.triple_history_size_after_last_prune, 10_000
):
return triple_history
cutoff_date = (
datetime.strptime(latest_event_date, "%Y-%m-%dT%H:%M:%SZ")
- timedelta(minutes=history_minutes)
).strftime("%Y-%m-%dT%H:%M:%SZ")
triple_history = {
triple: history
for triple, history in triple_history.items()
if history[2] >= cutoff_date
}
self.triple_history_size_after_last_prune = len(triple_history)
return triple_history

def check_and_update_triple_history(
self,
triple_history: dict[str, tuple[tuple[int, int], bool, str]],
triple: str,
rev_key: tuple[int, int],
is_insert: bool,
event_date: str,
) -> bool:
"""
Check the given event against the history of the causally newest
event per triple (see `triple_history` in `execute`) and return
whether the event should be applied. An insert is superseded by a
causally strictly newer delete of the same triple, a delete by a
causally newer or equally new insert (an insert wins a tie, a
delete does not, matching the per-batch tie semantics). If the
event is the causally newest one for its triple so far, it is
recorded in the history (which is modified in place).
"""
history = triple_history.get(triple)
is_newest = history is None or (
rev_key >= history[0] if is_insert else rev_key > history[0]
)
if is_newest:
triple_history[triple] = (rev_key, is_insert, event_date)
return True

# An event that is not the causally newest one is superseded only
# by an event of the other kind (a delete by an insert and vice
# versa); a duplicate of the same kind is harmless.
return history[1] == is_insert

def execute(self, args) -> bool:
# Resolve the four args that depend on `--wikimedia-commons`. A
# `None` value means the user did not pass that arg, so fill in
Expand Down Expand Up @@ -512,6 +592,24 @@ def execute(self, args) -> bool:
log.error(f"Error determining offset from stream: {e}")
return False

# History of the causal-order key of the chronologically newest
# event seen for each triple, kept ACROSS batches: maps the triple
# to `(rev_key, is_insert, event_date)`. The per-batch
# `insert_triples` / `delete_triples` below catch out-of-order
# events WITHIN a batch, but the stream can reorder events by
# minutes, while batches span only seconds when tailing the live
# stream. Without this history, a delete that arrives after the
# causally newer insert of the same triple was applied in an
# earlier batch deletes that triple (observed live: a sitelink
# moved from a duplicate entity to the right one, where the
# delete for the duplicate arrived two minutes late and stripped
# the article triples of the right entity). Entries older than
# `--triple-history-minutes` of stream time are pruned at each
# batch start.
triple_history: dict[str, tuple[tuple[int, int], bool, str]] = {}
latest_event_date = None
use_triple_history = args.triple_history_minutes > 0

# Initialize all the statistics variables.
batch_count = 0
total_num_messages = 0
Expand Down Expand Up @@ -695,12 +793,27 @@ def execute(self, args) -> bool:
# DELETE and the target's ADD for a transferred article URL share
# the same `dt` and the target often arrives first). Tracking the
# `rev_id` per triple makes the cross-set "remove from the other
# side" step gated by causal order, so the chronologically last
# event for a given triple wins regardless of arrival order.
# side" step respect the causal order, so the chronologically
# last event for a given triple wins regardless of arrival order.
insert_triples: dict[str, tuple[int, int]] = {}
delete_triples: dict[str, tuple[int, int]] = {}

# Prune old entries from the triple history.
if use_triple_history:
triple_history = self.prune_triple_history(
triple_history,
latest_event_date,
args.triple_history_minutes,
)

# Check if we can use a cached SPARQL query file
Comment on lines +805 to 809
#
# NOTE: A cached batch is applied without processing its events,
# so it contributes nothing to `triple_history`. For events in
# such a batch, the protection against out-of-order events is
# therefore only as good as it was before the introduction of
# `triple_history` (namely, within the batch that produced the
# cached file).
use_cached_file = False
cached_file_name = None
if (
Expand Down Expand Up @@ -746,6 +859,11 @@ def execute(self, args) -> bool:
# Get the date (rounded *down* to seconds).
date = meta.get("dt")
date = re.sub(r"\.\d*Z$", "Z", date)
if (
latest_event_date is None
or date > latest_event_date
):
latest_event_date = date

# Get the other relevant fields from the message.
entity_id = event_data.get("entity_id")
Expand Down Expand Up @@ -892,6 +1010,19 @@ def node_to_sparql(node: rdflib.term.Node) -> str:
)
for s, p, o in graph:
triple = f"{s.n3()} {p.n3()} {node_to_sparql(o)}"
# Cross-batch check, see
# `check_and_update_triple_history`.
if (
use_triple_history
and not self.check_and_update_triple_history(
triple_history,
triple,
rev_key,
is_insert=False,
event_date=date,
)
):
continue
# NOTE: In case there was a previous `insert` of that
# triple, it is safe to remove that `insert`, but not
# the `delete` (in case the triple is contained in the
Expand Down Expand Up @@ -942,10 +1073,23 @@ def node_to_sparql(node: rdflib.term.Node) -> str:
)
for s, p, o in graph:
triple = f"{s.n3()} {p.n3()} {node_to_sparql(o)}"
# Cross-batch check, see
# `check_and_update_triple_history`.
if (
use_triple_history
and not self.check_and_update_triple_history(
triple_history,
triple,
rev_key,
is_insert=True,
event_date=date,
)
):
continue
# NOTE: In case there was a previous `delete` of that
# triple, it is safe to remove that `delete`, but not
# the `insert` (in case the triple is not contained in
# the original data). Use `>=` on the cross-set gate
# the original data). Use `>=` on the cross-set check
# so that a same-event delete-then-add (deletes are
# processed first in the per-event loop) lets the
# add win, matching the previous within-event
Expand Down Expand Up @@ -980,6 +1124,17 @@ def node_to_sparql(node: rdflib.term.Node) -> str:

# Message was successfully processed, update batch tracking
current_batch_size += 1
if (
use_triple_history
and current_batch_size % 10000 == 0
):
# A batch can span hours of stream time during
# catch-up, so also prune within a batch.
triple_history = self.prune_triple_history(
triple_history,
latest_event_date,
args.triple_history_minutes,
)
total_num_messages += 1
pbar_update_frequency = 100
if (current_batch_size % pbar_update_frequency) == 0:
Expand Down
Loading