aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--generalresearch/incite/collections/base.py84
-rw-r--r--generalresearch/incite/mergers/pop_ledger.py41
-rw-r--r--generalresearch/incite/schemas/mergers/pop_ledger.py6
-rw-r--r--generalresearch/incite/schemas/thl_web.py10
-rw-r--r--pyproject.toml2
5 files changed, 82 insertions, 61 deletions
diff --git a/generalresearch/incite/collections/base.py b/generalresearch/incite/collections/base.py
index 04dd97e..f8159e8 100644
--- a/generalresearch/incite/collections/base.py
+++ b/generalresearch/incite/collections/base.py
@@ -1,6 +1,5 @@
from __future__ import annotations
-import logging
import os
import subprocess
import time
@@ -20,7 +19,6 @@ from dask.distributed import Future
from distributed import as_completed
from more_itertools import chunked
from pandera.pandas import DataFrameSchema
-from psycopg import Cursor
from pydantic import Field, FilePath, ValidationInfo, field_validator
from sentry_sdk import capture_exception
@@ -313,34 +311,36 @@ class DFCollectionItem(CollectionItemBase):
coll = self._collection
pg_config: PostgresConfig = coll.pg_config
- limit = 20000
+ limit = 10000
offset = 0
res = []
+ query_base = """
+ SELECT lt.id AS tx_id, lt.created, lt.ext_description, lt.tag,
+ le.id AS entry_id, le.direction, le.amount, le.account_id,
+ la.display_name, la.qualified_name, la.account_type,
+ la.normal_balance, la.reference_type, la.reference_uuid,
+ la.currency, tu.product_id, tu.product_user_id
+ FROM ledger_transaction AS lt
+ LEFT JOIN ledger_entry AS le
+ ON lt.id = le.transaction_id
+ LEFT JOIN ledger_account AS la
+ ON la.uuid = le.account_id
+ LEFT JOIN thl_user AS tu
+ ON tu.uuid = la.reference_uuid AND la.reference_type = 'user'
+ WHERE lt.created >= %s AND lt.created < %s
+ AND le.id IS NOT NULL
+ ORDER BY lt.created
+ """
while True:
- logging.info(f"{data_type.value}.from_postgres_ledger({limit=}, {offset=})")
- chunk = pg_config.execute_sql_query(
- query=f"""
- SELECT lt.id AS tx_id, lt.created, lt.ext_description, lt.tag,
- le.id AS entry_id, le.direction, le.amount, le.account_id,
- la.display_name, la.qualified_name, la.account_type,
- la.normal_balance, la.reference_type, la.reference_uuid,
- la.currency
- FROM ledger_transaction AS lt
- LEFT JOIN ledger_entry AS le
- ON lt.id = le.transaction_id
- LEFT JOIN ledger_account AS la
- ON la.uuid = le.account_id
- WHERE lt.created >= %s AND lt.created < %s
- AND le.id IS NOT NULL
- ORDER BY lt.created
- LIMIT {limit} OFFSET {offset};
- """,
- params=[start, finish],
- )
- res.extend(chunk)
- if not chunk:
- break
- offset += limit
+ with pg_config.make_connection() as conn, conn.cursor() as c:
+ query = query_base + f"\nLIMIT {limit} OFFSET {offset};"
+ c.execute(query, params=[start, finish])
+ chunk = c.fetchall()
+ LOG.info(f"{data_type.value}.from_postgres_ledger({limit=}, {offset=})")
+ res.extend(chunk)
+ if not chunk or len(chunk) < limit:
+ break
+ offset += len(chunk)
if len(res) == 0:
return None
@@ -357,23 +357,19 @@ class DFCollectionItem(CollectionItemBase):
tx_ids = list(tx_df["tx_id"].unique())
metadata_res = []
- # "MySQL server has gone away" if this is too big
- conn = pg_config.make_connection()
- c: Cursor = conn.cursor()
- for chunk in chunked(tx_ids, n=5_000):
- c.execute(
- query="""
- SELECT ltm.transaction_id AS tx_id,
- ltm.id AS tx_metadata_id,
- ltm.key, ltm.value
- FROM ledger_transactionmetadata AS ltm
- WHERE ltm.transaction_id = ANY(%s);
- """,
- params=[chunk],
- )
- metadata_res += c.fetchall()
-
- conn.close()
+ with pg_config.make_connection() as conn, conn.cursor() as c:
+ for chunk in chunked(tx_ids, n=5_000):
+ c.execute(
+ query="""
+ SELECT ltm.transaction_id AS tx_id,
+ ltm.id AS tx_metadata_id,
+ ltm.key, ltm.value
+ FROM ledger_transactionmetadata AS ltm
+ WHERE ltm.transaction_id = ANY(%s);
+ """,
+ params=[chunk],
+ )
+ metadata_res += c.fetchall()
tx_meta = (
pd.DataFrame(
diff --git a/generalresearch/incite/mergers/pop_ledger.py b/generalresearch/incite/mergers/pop_ledger.py
index 54f1b7e..9120ca0 100644
--- a/generalresearch/incite/mergers/pop_ledger.py
+++ b/generalresearch/incite/mergers/pop_ledger.py
@@ -19,7 +19,6 @@ from generalresearch.models.thl.ledger import Direction, TransactionType
class PopLedgerMergeItem(MergeCollectionItem):
-
def build(
self,
ledger_coll: LedgerDFCollection,
@@ -63,24 +62,35 @@ class PopLedgerMergeItem(MergeCollectionItem):
# For each time interval and Ledger Account (this is different from a
# product_id), we want the raw amounts and their respective direction
# for every type of transaction that is possible
- x = (
- df.groupby(by=["time_idx", "account_id", "tx_type", "direction_name"])
- .amount.sum()
- .reset_index()
- )
+ gb_cols = [
+ "time_idx",
+ "account_id",
+ "tx_type",
+ "direction_name",
+ "reference_uuid",
+ "product_id",
+ "product_user_id",
+ ]
+ x = df.groupby(by=gb_cols, dropna=False).amount.sum().reset_index()
# We want to keep the positive and negatives for each type. For example,
# for bp_adjustment, we want to know the amount increased and the
# amount decreased, not just the net.
x["tx_type.direction"] = x["tx_type"] + "." + x["direction_name"]
+ identity_cols = [
+ "time_idx",
+ "account_id",
+ "product_id",
+ "product_user_id",
+ ]
s = (
- x.pivot_table(
- index=["time_idx", "account_id"],
- columns="tx_type.direction",
- values="amount",
- aggfunc="sum",
- )
- .fillna(0)
+ x.groupby(
+ identity_cols + ["tx_type.direction"],
+ dropna=False,
+ observed=True,
+ )["amount"]
+ .sum()
+ .unstack(fill_value=0)
.reset_index()
)
@@ -89,7 +99,10 @@ class PopLedgerMergeItem(MergeCollectionItem):
[[e.value + ".CREDIT", e.value + ".DEBIT"] for e in TransactionType]
)
)
- s = s.reindex(columns=columns | set(s.columns)).fillna(0)
+ for column in columns.difference(s.columns):
+ s[column] = 0
+ s[list(columns)] = s[list(columns)].fillna(0)
+
s = s.reset_index(drop=True)
s.index.name = "id"
# The "columns were named" tx_type.direction from the above pivot. This
diff --git a/generalresearch/incite/schemas/mergers/pop_ledger.py b/generalresearch/incite/schemas/mergers/pop_ledger.py
index 8452eb9..eb81648 100644
--- a/generalresearch/incite/schemas/mergers/pop_ledger.py
+++ b/generalresearch/incite/schemas/mergers/pop_ledger.py
@@ -25,8 +25,8 @@ from generalresearch.models.thl.ledger import Direction, TransactionType
# If an amount is "very" large, something is def wrong. Defining "very" somewhat arbitrarily here.
SUSPICIOUSLY_LARGE_NUMBER = (2**32 / 2) - 1 # 2147483647
-_tz_min_freq: Callable[[pd.Series], pd.Series] = lambda i: (i.dt.second == 0) & (
- i.dt.microsecond == 0
+_tz_min_freq: Callable[[pd.Series], pd.Series] = lambda i: (
+ (i.dt.second == 0) & (i.dt.microsecond == 0)
)
@@ -61,6 +61,8 @@ PopLedgerSchema = DataFrameSchema(
nullable=False,
),
"account_id": TxSchema.columns["account_id"],
+ "product_id": TxSchema.columns["product_id"],
+ "product_user_id": TxSchema.columns["product_user_id"],
},
checks=[],
coerce=True,
diff --git a/generalresearch/incite/schemas/thl_web.py b/generalresearch/incite/schemas/thl_web.py
index 36ee8e9..b816b9f 100644
--- a/generalresearch/incite/schemas/thl_web.py
+++ b/generalresearch/incite/schemas/thl_web.py
@@ -739,6 +739,16 @@ TxSchema = DataFrameSchema(
],
nullable=False,
),
+ "product_id": Column(
+ dtype=str,
+ checks=Check.str_length(min_value=32, max_value=32),
+ nullable=True,
+ ),
+ "product_user_id": Column(
+ dtype=str,
+ checks=Check.str_length(min_value=3, max_value=128),
+ nullable=True,
+ ),
},
checks=[],
coerce=True,
diff --git a/pyproject.toml b/pyproject.toml
index 19619aa..76aa43a 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "generalresearch"
-version = "3.6.5"
+version = "3.6.6"
description = "Python Utilities for General Research"
readme = "README.md"
requires-python = ">=3.14"