Hi All , I was working on a chain reader that retrieves a chain and all subchains down to RIC level. We already have a chain dissector, but it is slow when a chain contains around 30K chains/subchains.
I also built a multithreaded chain reader, but the LSEG API starts blocking when more than two worker threads are used. If anyone wants to build on this, or can implement an end-to-end version that returns full chain details in about 20 minutes, please let me know if you succeed.
Example chain: 0#ARGUSMEDIA.
import csv
import logging
import os
import threading
import time
import warnings
from collections import deque
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime
import pandas as pd
import lseg.data as ld
logging.getLogger("lseg").setLevel(logging.ERROR)
warnings.simplefilter("ignore", FutureWarning)
pd.set_option("future.no_silent_downcasting", True)
LONGNEXT_FIELD = "LONGNEXTLR"
LONGLINK_FIELDS = [f"LONGLINK{i}" for i in range(1, 15)]
ALL_FIELDS = [LONGNEXT_FIELD] + LONGLINK_FIELDS
REQUEST_PAUSE_SEC = 0.02
BATCH_SIZE = 10
INITIAL_WORKERS = int(os.getenv("CHAIN_WORKERS", "2"))
MAX_WORKERS = int(os.getenv("CHAIN_MAX_WORKERS", str(max(2, INITIAL_WORKERS))))
ALERT_EVERY_CHAINS = 100
MAX_BATCH_RETRIES = 3
RECONNECT_BACKOFF_BASE_SEC = 2.0
THREAD_START_STAGGER_SEC = 1.0
RAMP_UP_AFTER_BRANCHES = 2
APP_KEY = "YOURAPPID"
USERNAME = "YOUR_UN"
PASSWORD = "YOUR_PW"
ROOT_RIC = "0#ARGUSMEDIA"
thread_local = threading.local()
def log(msg: str) -> None:
print(msg, flush=True)
def _new_session():
session = ld.session.platform.Definition(
app_key=APP_KEY,
grant=ld.session.platform.GrantPassword(
username=USERNAME,
password=PASSWORD,
),
signon_control=True,
).get_session()
session.open()
if "Opened" not in str(session.open_state):
raise RuntimeError(f"Session failed to open: {session.open_state}")
ld.session.set_default(session)
return session
def _get_thread_session(force_reopen: bool = False):
session = getattr(thread_local, "session", None)
if session is not None and not force_reopen:
try:
if "Opened" in str(session.open_state):
ld.session.set_default(session)
return session
except Exception:
pass
if session is not None:
try:
session.close()
except Exception:
pass
last_err = None
for attempt in range(5):
try:
session = _new_session()
thread_local.session = session
log(f"[session] {threading.current_thread().name} connected")
return session
except Exception as exc:
last_err = exc
wait_sec = min(30, 2 ** attempt)
log(
f"[session] {threading.current_thread().name} open failed "
f"({attempt + 1}/5): {exc}; retry in {wait_sec}s"
)
time.sleep(wait_sec)
raise RuntimeError(f"Unable to open session: {last_err}")
def _close_thread_session() -> None:
session = getattr(thread_local, "session", None)
if session is not None:
try:
session.close()
except Exception:
pass
thread_local.session = None
def _fetch_rows(universe: list[str], max_retries: int = MAX_BATCH_RETRIES) -> dict[str, dict]:
if not universe:
return {}
wanted = set(universe)
last_err = None
for attempt in range(max_retries):
try:
_get_thread_session()
if REQUEST_PAUSE_SEC:
time.sleep(REQUEST_PAUSE_SEC)
df = ld.get_data(universe=universe, fields=ALL_FIELDS)
if df is None or df.empty:
return {ric: {} for ric in universe}
id_col = None
for candidate in ("Instrument", "instrument", "RIC", "ric", "Universe", "universe"):
if candidate in df.columns:
id_col = candidate
break
if id_col is None:
id_col = df.columns[0]
out = {}
for _, row in df.iterrows():
key = str(row.get(id_col, "")).strip()
if key in wanted:
out[key] = dict(row)
for ric in universe:
out.setdefault(ric, {})
return out
except Exception as exc:
last_err = exc
err_str = str(exc).lower()
if any(k in err_str for k in ("lost", "closed", "socket", "connection", "stream", "timeout")):
log(f"[retry] {threading.current_thread().name} reconnect for batch: {exc}")
try:
_get_thread_session(force_reopen=True)
except Exception as re:
log(f"[retry] {threading.current_thread().name} reconnect failed: {re}")
time.sleep(min(15.0, RECONNECT_BACKOFF_BASE_SEC * (attempt + 1)))
else:
time.sleep(0.25 * (attempt + 1))
log(f"[skip] batch failed after retries: {last_err}")
return {ric: {} for ric in universe}
def _normalize_values(row: dict) -> list[str]:
values = []
for field in LONGLINK_FIELDS:
val = row.get(field)
if val is None:
continue
text = str(val).strip()
if text and text.lower() not in ("nan", "none", ""):
values.append(text)
return values
def _expand_chain_pages(chain_ric: str, first_row: dict | None = None) -> list[str]:
members = []
seen_members = set()
seen_pages = set()
current = chain_ric
row = first_row
while current and current not in seen_pages:
seen_pages.add(current)
if row is None:
row = _fetch_rows([current]).get(current, {})
if not row:
break
for symbol in _normalize_values(row):
if symbol not in seen_members:
seen_members.add(symbol)
members.append(symbol)
next_chain = row.get(LONGNEXT_FIELD)
if next_chain is None:
break
next_text = str(next_chain).strip()
if not next_text or next_text.lower() in ("nan", "none", ""):
break
current = next_text
row = None
return members
def _build_branch_tree(branch_root: str) -> dict[str, dict[str, list[str]]]:
tree = {}
visited = {branch_root}
queue = deque([branch_root])
done = 0
started = time.time()
# Stagger worker startup to reduce simultaneous stream/session initialization pressure.
name = threading.current_thread().name
if name.startswith("branch_"):
try:
idx = int(name.split("_")[-1])
time.sleep(idx * THREAD_START_STAGGER_SEC)
except Exception:
pass
log(f"[branch-start] {branch_root}")
while queue:
batch = []
while queue and len(batch) < BATCH_SIZE:
batch.append(queue.popleft())
first_rows = _fetch_rows(batch)
for chain_ric in batch:
members = _expand_chain_pages(chain_ric, first_row=first_rows.get(chain_ric))
subchains = [m for m in members if "#" in m]
rics = [m for m in members if "#" not in m]
tree[chain_ric] = {"rics": rics, "subchains": subchains}
done += 1
for sub in subchains:
if sub not in visited:
visited.add(sub)
queue.append(sub)
if done % ALERT_EVERY_CHAINS == 0:
elapsed = time.time() - started
log(
f" [branch] {branch_root} done={done} queued={len(queue)} elapsed={elapsed:.0f}s"
)
elapsed = time.time() - started
log(f"[branch-done] {branch_root} chains={done} elapsed={elapsed:.0f}s")
_close_thread_session()
return tree
def write_csv(tree: dict[str, dict[str, list[str]]], root: str, output_path: str) -> None:
rows = []
seen = set()
def walk(chain_ric: str, level: int = 0) -> None:
if chain_ric in seen:
return
seen.add(chain_ric)
node = tree.get(chain_ric, {"rics": [], "subchains": []})
rows.append(["CHAIN", chain_ric, f"LEVEL_{level}"])
rows.append(["TYPE", "VALUE", ""])
for sub in node.get("subchains", []):
rows.append(["subchain", sub, ""])
for ric in node.get("rics", []):
rows.append(["ric", ric, ""])
rows.append([])
for sub in node.get("subchains", []):
walk(sub, level + 1)
walk(root)
with open(output_path, "w", newline="", encoding="utf-8") as f:
csv.writer(f).writerows(rows)
log(f"[csv] written: {output_path}")
def main() -> None:
overall_start = time.time()
start_workers = min(INITIAL_WORKERS, MAX_WORKERS)
max_workers = max(start_workers, MAX_WORKERS)
log(f"[start] root={ROOT_RIC} workers={start_workers} max_workers={max_workers}")
_get_thread_session()
root_members = _expand_chain_pages(ROOT_RIC)
root_subchains = [m for m in root_members if "#" in m]
root_rics = [m for m in root_members if "#" not in m]
root_tree = {ROOT_RIC: {"rics": root_rics, "subchains": root_subchains}}
log(
f"[root] subchains={len(root_subchains)} rics={len(root_rics)} "
f"top-level branches discovered"
)
_close_thread_session()
merged_tree = dict(root_tree)
with ThreadPoolExecutor(max_workers=max_workers, thread_name_prefix="branch") as executor:
pending = deque(root_subchains)
future_map: dict = {}
active_target = start_workers
completed_ok = 0
def submit_until_full() -> None:
while pending and len(future_map) < active_target:
branch_root = pending.popleft()
future = executor.submit(_build_branch_tree, branch_root)
future_map[future] = branch_root
submit_until_full()
while future_map:
for future in as_completed(list(future_map.keys()), timeout=None):
branch_root = future_map.pop(future)
try:
branch_tree = future.result()
merged_tree.update(branch_tree)
completed_ok += 1
except Exception as exc:
log(f"[branch-error] {branch_root}: {exc}")
if pending and active_target < max_workers and completed_ok % RAMP_UP_AFTER_BRANCHES == 0:
active_target += 1
log(f"[ramp] workers={active_target} pending={len(pending)}")
submit_until_full()
break
total_chains = len(merged_tree)
total_leaf_rics = sum(len(v.get("rics", [])) for v in merged_tree.values())
unique_leaf_rics = len({ric for v in merged_tree.values() for ric in v.get("rics", [])})
elapsed = time.time() - overall_start
log(f"[done] chains={total_chains} leaf_refs={total_leaf_rics} unique_leaf_rics={unique_leaf_rics} elapsed={elapsed:.0f}s")
stamp = datetime.now().strftime("%Y%m%d_%H%M%S")
output_path = os.path.join(os.getcwd(), f"chain_tree_fast_{stamp}.csv")
write_csv(merged_tree, ROOT_RIC, output_path)
if name == "main":
main()