Skip to content

Commit dae3ec9

Browse files
authored
Merge pull request #225 from posit-dev/fix-centralize-data-ingest
fix: use `_process_data()` to centralize data ingest functionality
2 parents 75a014c + 168aabe commit dae3ec9

6 files changed

Lines changed: 332 additions & 156 deletions

File tree

pointblank/assistant.py

Lines changed: 2 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -180,21 +180,10 @@ def assistant(
180180
# If a dataset is provided, generate a table summary in JSON format
181181
if data is not None:
182182
# Import processing functions from validate module
183-
from pointblank.validate import (
184-
_process_connection_string,
185-
_process_csv_input,
186-
_process_parquet_input,
187-
)
183+
from pointblank.validate import _process_data
188184

189185
# Process input data to handle different data source types
190-
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
191-
data = _process_connection_string(data)
192-
193-
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
194-
data = _process_csv_input(data)
195-
196-
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
197-
data = _process_parquet_input(data)
186+
data = _process_data(data)
198187

199188
scan = DataScan(data=data)
200189

pointblank/cli.py

Lines changed: 12 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -832,18 +832,12 @@ def preview(
832832

833833
# If data has _row_num_ and it's not explicitly included, add it at the beginning
834834
try:
835-
from pointblank.validate import (
836-
_process_connection_string,
837-
_process_csv_input,
838-
_process_parquet_input,
839-
)
835+
from pointblank.validate import _process_data
840836

841837
# Process the data source to get actual data object to check for _row_num_
842838
processed_data = data
843839
if isinstance(data, str):
844-
processed_data = _process_connection_string(data)
845-
processed_data = _process_csv_input(processed_data)
846-
processed_data = _process_parquet_input(processed_data)
840+
processed_data = _process_data(data)
847841

848842
# Get column names from the processed data
849843
all_columns = []
@@ -861,18 +855,12 @@ def preview(
861855
elif col_range or col_first or col_last:
862856
# Need to get column names to apply range/first/last selection
863857
# Load the data to get column names
864-
from pointblank.validate import (
865-
_process_connection_string,
866-
_process_csv_input,
867-
_process_parquet_input,
868-
)
858+
from pointblank.validate import _process_data
869859

870860
# Process the data source to get actual data object
871861
processed_data = data
872862
if isinstance(data, str):
873-
processed_data = _process_connection_string(data)
874-
processed_data = _process_csv_input(processed_data)
875-
processed_data = _process_parquet_input(processed_data)
863+
processed_data = _process_data(data)
876864

877865
# Get column names from the processed data
878866
all_columns = []
@@ -935,17 +923,11 @@ def preview(
935923
# Get total dataset size before preview and gather metadata
936924
try:
937925
# Process the data to get the actual data object for row count and metadata
938-
from pointblank.validate import (
939-
_process_connection_string,
940-
_process_csv_input,
941-
_process_parquet_input,
942-
)
926+
from pointblank.validate import _process_data
943927

944928
processed_data = data
945929
if isinstance(data, str):
946-
processed_data = _process_connection_string(data)
947-
processed_data = _process_csv_input(processed_data)
948-
processed_data = _process_parquet_input(processed_data)
930+
processed_data = _process_data(data)
949931

950932
total_dataset_rows = pb.get_row_count(processed_data)
951933

@@ -1024,15 +1006,9 @@ def info(data_source: str):
10241006
source_type = f"External source: {data_source}"
10251007

10261008
# Process the data to get actual table object for inspection
1027-
from pointblank.validate import (
1028-
_process_connection_string,
1029-
_process_csv_input,
1030-
_process_parquet_input,
1031-
)
1009+
from pointblank.validate import _process_data
10321010

1033-
data = _process_connection_string(data)
1034-
data = _process_csv_input(data)
1035-
data = _process_parquet_input(data)
1011+
data = _process_data(data)
10361012
console.print(f"[green]✓[/green] Loaded data source: {data_source}")
10371013

10381014
# Get table information
@@ -1131,15 +1107,9 @@ def scan(
11311107
total_rows = None
11321108
else:
11331109
# For file paths and connection strings, load the data first
1134-
from pointblank.validate import (
1135-
_process_connection_string,
1136-
_process_csv_input,
1137-
_process_parquet_input,
1138-
)
1110+
from pointblank.validate import _process_data
11391111

1140-
processed_data = _process_connection_string(data)
1141-
processed_data = _process_csv_input(processed_data)
1142-
processed_data = _process_parquet_input(processed_data)
1112+
processed_data = _process_data(data)
11431113
scan_result = pb.col_summary_tbl(data=processed_data)
11441114
source_type = f"External source: {data_source}"
11451115
table_type = _get_tbl_type(processed_data)
@@ -1212,16 +1182,10 @@ def missing(data_source: str, output_html: str | None):
12121182
original_data = data
12131183
if isinstance(data, str):
12141184
# Process the data to get the actual data object
1215-
from pointblank.validate import (
1216-
_process_connection_string,
1217-
_process_csv_input,
1218-
_process_parquet_input,
1219-
)
1185+
from pointblank.validate import _process_data
12201186

12211187
try:
1222-
original_data = _process_connection_string(data)
1223-
original_data = _process_csv_input(original_data)
1224-
original_data = _process_parquet_input(original_data)
1188+
original_data = _process_data(data)
12251189
except Exception: # pragma: no cover
12261190
pass # Use the string data as fallback
12271191

pointblank/compare.py

Lines changed: 3 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -11,21 +11,13 @@
1111
class Compare:
1212
def __init__(self, a: IntoFrame, b: IntoFrame) -> None:
1313
# Import processing functions from validate module
14-
from pointblank.validate import (
15-
_process_connection_string,
16-
_process_csv_input,
17-
_process_parquet_input,
18-
)
14+
from pointblank.validate import _process_data
1915

2016
# Process input data for table a
21-
a = _process_connection_string(a)
22-
a = _process_csv_input(a)
23-
a = _process_parquet_input(a)
17+
a = _process_data(a)
2418

2519
# Process input data for table b
26-
b = _process_connection_string(b)
27-
b = _process_csv_input(b)
28-
b = _process_parquet_input(b)
20+
b = _process_data(b)
2921

3022
self.a: IntoFrame = a
3123
self.b: IntoFrame = b

pointblank/datascan.py

Lines changed: 4 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -126,24 +126,11 @@ class DataScan:
126126
def __init__(self, data: IntoFrameT, tbl_name: str | None = None) -> None:
127127
# Import processing functions from validate module
128128
from pointblank.validate import (
129-
_process_connection_string,
130-
_process_csv_input,
131-
_process_github_url,
132-
_process_parquet_input,
129+
_process_data,
133130
)
134131

135132
# Process input data to handle different data source types
136-
# Handle GitHub URL input (e.g., "https://github.com/user/repo/blob/main/data.csv")
137-
data = _process_github_url(data)
138-
139-
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
140-
data = _process_connection_string(data)
141-
142-
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
143-
data = _process_csv_input(data)
144-
145-
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
146-
data = _process_parquet_input(data)
133+
data = _process_data(data)
147134

148135
as_native = nw.from_native(data)
149136

@@ -596,25 +583,10 @@ def col_summary_tbl(data: FrameT | Any, tbl_name: str | None = None) -> GT:
596583
"""
597584

598585
# Import processing functions from validate module
599-
from pointblank.validate import (
600-
_process_connection_string,
601-
_process_csv_input,
602-
_process_github_url,
603-
_process_parquet_input,
604-
)
586+
from pointblank.validate import _process_data
605587

606588
# Process input data to handle different data source types
607-
# Handle GitHub URL input (e.g., "https://github.com/user/repo/blob/main/data.csv")
608-
data = _process_github_url(data)
609-
610-
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
611-
data = _process_connection_string(data)
612-
613-
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
614-
data = _process_csv_input(data)
615-
616-
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
617-
data = _process_parquet_input(data)
589+
data = _process_data(data)
618590

619591
scanner = DataScan(data=data, tbl_name=tbl_name)
620592
return scanner.get_tabular_report()

pointblank/validate.py

Lines changed: 51 additions & 52 deletions
Original file line numberDiff line numberDiff line change
@@ -735,9 +735,51 @@ def get_data_path(
735735
return tmp_file.name
736736

737737

738-
# =============================================================================
739-
# Utility functions for processing input data (shared by preview() and Validate class)
740-
# =============================================================================
738+
def _process_data(data: FrameT | Any) -> FrameT | Any:
739+
"""
740+
Centralized data processing pipeline that handles all supported input types.
741+
742+
This function consolidates the data processing pipeline used across multiple
743+
classes and functions in Pointblank. It processes data through a consistent
744+
sequence of transformations to handle different data source types.
745+
746+
The processing order is important:
747+
748+
1. GitHub URLs (must come before connection string processing)
749+
2. Database connection strings
750+
3. CSV file paths
751+
4. Parquet file paths
752+
753+
Parameters
754+
----------
755+
data : FrameT | Any
756+
The input data which could be:
757+
- a DataFrame object (Polars, Pandas, Ibis, etc.)
758+
- a GitHub URL pointing to a CSV or Parquet file
759+
- a database connection string (e.g., "duckdb:///path/to/file.ddb::table_name")
760+
- a CSV file path (string or Path object with .csv extension)
761+
- a Parquet file path, glob pattern, directory, or partitioned dataset
762+
- any other data type (returned unchanged)
763+
764+
Returns
765+
-------
766+
FrameT | Any
767+
Processed data as a DataFrame if input was a supported data source type,
768+
otherwise the original data unchanged.
769+
"""
770+
# Handle GitHub URL input (e.g., "https://github.com/user/repo/blob/main/data.csv")
771+
data = _process_github_url(data)
772+
773+
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
774+
data = _process_connection_string(data)
775+
776+
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
777+
data = _process_csv_input(data)
778+
779+
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
780+
data = _process_parquet_input(data)
781+
782+
return data
741783

742784

743785
def _process_github_url(data: FrameT | Any) -> FrameT | Any:
@@ -1321,17 +1363,7 @@ def preview(
13211363
"""
13221364

13231365
# Process input data to handle different data source types
1324-
# Handle GitHub URL input (e.g., "https://github.com/user/repo/blob/main/data.csv")
1325-
data = _process_github_url(data)
1326-
1327-
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
1328-
data = _process_connection_string(data)
1329-
1330-
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
1331-
data = _process_csv_input(data)
1332-
1333-
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
1334-
data = _process_parquet_input(data)
1366+
data = _process_data(data)
13351367

13361368
if incl_header is None:
13371369
incl_header = global_config.preview_incl_header
@@ -1816,17 +1848,7 @@ def missing_vals_tbl(data: FrameT | Any) -> GT:
18161848
"""
18171849

18181850
# Process input data to handle different data source types
1819-
# Handle GitHub URL input (e.g., "https://github.com/user/repo/blob/main/data.csv")
1820-
data = _process_github_url(data)
1821-
1822-
# Handle connection string input (e.g., "duckdb:///path/to/file.ddb::table_name")
1823-
data = _process_connection_string(data)
1824-
1825-
# Handle CSV file input (e.g., "data.csv" or Path("data.csv"))
1826-
data = _process_csv_input(data)
1827-
1828-
# Handle Parquet file input (e.g., "data.parquet", "data/*.parquet", "data/")
1829-
data = _process_parquet_input(data)
1851+
data = _process_data(data)
18301852

18311853
# Make a copy of the data to avoid modifying the original
18321854
data = copy.deepcopy(data)
@@ -2431,14 +2453,7 @@ def get_column_count(data: FrameT | Any) -> int:
24312453

24322454
# Process different input types
24332455
if isinstance(data, str) or isinstance(data, Path):
2434-
# Process GitHub URLs first
2435-
data = _process_github_url(data)
2436-
# Handle connection string input
2437-
data = _process_connection_string(data)
2438-
# Handle CSV file input
2439-
data = _process_csv_input(data)
2440-
# Handle Parquet file input
2441-
data = _process_parquet_input(data)
2456+
data = _process_data(data)
24422457
elif isinstance(data, list):
24432458
# Handle list of file paths (likely Parquet files)
24442459
data = _process_parquet_input(data)
@@ -2607,14 +2622,7 @@ def get_row_count(data: FrameT | Any) -> int:
26072622

26082623
# Process different input types
26092624
if isinstance(data, str) or isinstance(data, Path):
2610-
# Process GitHub URLs first
2611-
data = _process_github_url(data)
2612-
# Handle connection string input
2613-
data = _process_connection_string(data)
2614-
# Handle CSV file input
2615-
data = _process_csv_input(data)
2616-
# Handle Parquet file input
2617-
data = _process_parquet_input(data)
2625+
data = _process_data(data)
26182626
elif isinstance(data, list):
26192627
# Handle list of file paths (likely Parquet files)
26202628
data = _process_parquet_input(data)
@@ -3550,17 +3558,8 @@ def send_report():
35503558
locale: str | None = None
35513559

35523560
def __post_init__(self):
3553-
# Handle GitHub URL input for the data parameter
3554-
self.data = _process_github_url(self.data)
3555-
3556-
# Handle connection string input for the data parameter
3557-
self.data = _process_connection_string(self.data)
3558-
3559-
# Handle CSV file input for the data parameter
3560-
self.data = _process_csv_input(self.data)
3561-
3562-
# Handle Parquet file input for the data parameter
3563-
self.data = _process_parquet_input(self.data)
3561+
# Process data through the centralized data processing pipeline
3562+
self.data = _process_data(self.data)
35643563

35653564
# Check input of the `thresholds=` argument
35663565
_check_thresholds(thresholds=self.thresholds)

0 commit comments

Comments
 (0)