diff options
| -rw-r--r-- | generalresearch/incite/collections/base.py | 84 | ||||
| -rw-r--r-- | generalresearch/incite/mergers/pop_ledger.py | 41 | ||||
| -rw-r--r-- | generalresearch/incite/schemas/mergers/pop_ledger.py | 6 | ||||
| -rw-r--r-- | generalresearch/incite/schemas/thl_web.py | 10 | ||||
| -rw-r--r-- | pyproject.toml | 2 |
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" |
