Skip to content

Commit a043eab

Browse files
committed
fix table identifier length
1 parent c9b7b99 commit a043eab

1 file changed

Lines changed: 54 additions & 28 deletions

File tree

main.py

Lines changed: 54 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -1,29 +1,44 @@
1+
import hashlib
12
import logging
2-
import pandas as pd
3-
from sqlalchemy import create_engine
43

5-
from typing import List
4+
import pandas as pd
5+
from preprocessing.core import Preprocessor
6+
from preprocessing.loader import BibLoader, TxtLoader
7+
from preprocessing.models import DocumentRecord, PreprocessedDocument
68
from scystream.sdk.core import entrypoint
79
from scystream.sdk.env.settings import (
810
EnvSettings,
11+
FileSettings,
912
InputSettings,
1013
OutputSettings,
11-
FileSettings,
12-
PostgresSettings
14+
PostgresSettings,
1315
)
1416
from scystream.sdk.file_handling.s3_manager import S3Operations
15-
16-
from preprocessing.core import Preprocessor
17-
from preprocessing.loader import TxtLoader, BibLoader
18-
from preprocessing.models import DocumentRecord, PreprocessedDocument
17+
from sqlalchemy import create_engine
18+
from sqlalchemy.sql import quoted_name
1919

2020
logging.basicConfig(
2121
level=logging.INFO,
22-
format='%(asctime)s - %(name)s - %(levelname)s - %(message)s'
22+
format="%(asctime)s - %(name)s - %(levelname)s - %(message)s",
2323
)
2424
logger = logging.getLogger(__name__)
2525

2626

27+
def _normalize_table_name(table_name: str) -> str:
28+
max_length = 63
29+
if len(table_name) <= max_length:
30+
return table_name
31+
digest = hashlib.sha1(table_name.encode("utf-8")).hexdigest()[:10]
32+
prefix_length = max_length - len(digest) - 1
33+
return f"{table_name[:prefix_length]}_{digest}"
34+
35+
36+
def _resolve_db_table(settings: PostgresSettings) -> str:
37+
normalized_name = _normalize_table_name(settings.DB_TABLE)
38+
settings.DB_TABLE = normalized_name
39+
return normalized_name
40+
41+
2742
class NormalizedDocsOutput(PostgresSettings, OutputSettings):
2843
__identifier__ = "normalized_docs"
2944

@@ -69,31 +84,40 @@ class PreprocessBIB(EnvSettings):
6984

7085

7186
def _write_preprocessed_docs_to_postgres(
72-
preprocessed_ouput: List[PreprocessedDocument],
73-
settings: PostgresSettings
87+
preprocessed_ouput: list[PreprocessedDocument],
88+
settings: PostgresSettings,
7489
):
75-
df = pd.DataFrame([
76-
{
77-
"doc_id": d.doc_id,
78-
"tokens": d.tokens
79-
}
80-
for d in preprocessed_ouput
81-
])
82-
83-
logger.info(f"Writing {len(df)} processed documents to DB table '{
84-
settings.DB_TABLE}'…")
90+
resolved_table_name = _resolve_db_table(settings)
91+
df = pd.DataFrame(
92+
[
93+
{
94+
"doc_id": d.doc_id,
95+
"tokens": d.tokens,
96+
}
97+
for d in preprocessed_ouput
98+
],
99+
)
100+
101+
logger.info(
102+
"Writing %s processed documents to DB table '%s'…",
103+
len(df),
104+
resolved_table_name,
105+
)
85106
engine = create_engine(
86107
f"postgresql+psycopg2://{settings.PG_USER}:{settings.PG_PASS}"
87-
f"@{settings.PG_HOST}:{int(settings.PG_PORT)}/"
108+
f"@{settings.PG_HOST}:{int(settings.PG_PORT)}/",
88109
)
89110

90-
df.to_sql(settings.DB_TABLE, engine, if_exists="replace", index=False)
111+
table_name = quoted_name(resolved_table_name, quote=True)
112+
df.to_sql(table_name, engine, if_exists="replace", index=False)
91113

92-
logger.info(f"Successfully stored normalized documents into '{
93-
settings.DB_TABLE}'.")
114+
logger.info(
115+
"Successfully stored normalized documents into '%s'.",
116+
resolved_table_name,
117+
)
94118

95119

96-
def _preprocess_and_store(documents: List[DocumentRecord], settings):
120+
def _preprocess_and_store(documents: list[DocumentRecord], settings):
97121
"""Shared preprocessing logic for TXT and BIB."""
98122
logger.info(f"Starting preprocessing with {len(documents)} documents")
99123

@@ -110,7 +134,9 @@ def _preprocess_and_store(documents: List[DocumentRecord], settings):
110134
result = pre.generate_normalized_output()
111135

112136
_write_preprocessed_docs_to_postgres(
113-
result, settings.normalized_docs_output)
137+
result,
138+
settings.normalized_docs_output,
139+
)
114140

115141
logger.info("Preprocessing completed successfully.")
116142

0 commit comments

Comments
 (0)