1717
1818import pandas as pd
1919
20+ from .._util import sanitize_csv_formulas
2021from ._cleaner import StreamingCleaner
2122
2223
@@ -34,10 +35,11 @@ def _read_chunks(path: str, batch_size: int) -> Iterator[pd.DataFrame]:
3435class _BatchWriter :
3536 """Append cleaned batches to one CSV or Parquet file without buffering them all."""
3637
37- def __init__ (self , path : str | None ) -> None :
38+ def __init__ (self , path : str | None , sanitize_formulas : bool = False ) -> None :
3839 self .path = path
3940 self .fmt = None if path is None else ("parquet"
4041 if path .lower ().endswith ((".parquet" , ".pq" )) else "csv" )
42+ self .sanitize_formulas = sanitize_formulas
4143 self ._pq_writer : Any = None
4244 self ._csv_header = True
4345
@@ -53,6 +55,8 @@ def write(self, df: pd.DataFrame) -> None:
5355 self ._pq_writer = pq .ParquetWriter (self .path , table .schema )
5456 self ._pq_writer .write_table (table )
5557 else :
58+ if self .sanitize_formulas :
59+ df = sanitize_csv_formulas (df )
5660 df .to_csv (self .path , mode = "w" if self ._csv_header else "a" ,
5761 header = self ._csv_header , index = False )
5862 self ._csv_header = False
@@ -109,7 +113,8 @@ def _timeseries_config(args: argparse.Namespace) -> Any:
109113 return TimeSeriesCleanConfig (** kwargs )
110114
111115
112- def _write_exceptions (cleaner : StreamingCleaner , path : str | None ) -> None :
116+ def _write_exceptions (cleaner : StreamingCleaner , path : str | None ,
117+ sanitize_formulas : bool = False ) -> None :
113118 """Persist any quarantined (late/anomalous) rows to *path* (CSV or Parquet)."""
114119 if not path :
115120 return
@@ -119,12 +124,15 @@ def _write_exceptions(cleaner: StreamingCleaner, path: str | None) -> None:
119124 if path .lower ().endswith ((".parquet" , ".pq" )):
120125 exc .to_parquet (path , index = False )
121126 else :
127+ if sanitize_formulas :
128+ exc = sanitize_csv_formulas (exc )
122129 exc .to_csv (path , index = False )
123130
124131
125132def _run_stream (cleaner : StreamingCleaner , batches : Iterator [pd .DataFrame ],
126133 writer : _BatchWriter , report_dir : str | None , quiet : bool ,
127- quarantine_path : str | None = None ) -> int :
134+ quarantine_path : str | None = None ,
135+ sanitize_formulas : bool = False ) -> int :
128136 if report_dir :
129137 os .makedirs (report_dir , exist_ok = True )
130138 for cleaned , report in cleaner .clean_batches (batches ):
@@ -137,7 +145,7 @@ def _run_stream(cleaner: StreamingCleaner, batches: Iterator[pd.DataFrame],
137145 print (json .dumps (report .streaming ))
138146 writer .close ()
139147 final = cleaner .finalize ()
140- _write_exceptions (cleaner , quarantine_path )
148+ _write_exceptions (cleaner , quarantine_path , sanitize_formulas = sanitize_formulas )
141149 if report_dir :
142150 with open (os .path .join (report_dir , "summary.json" ), "w" ) as fh :
143151 json .dump (final .to_dict (), fh , default = str )
@@ -148,17 +156,21 @@ def _run_stream(cleaner: StreamingCleaner, batches: Iterator[pd.DataFrame],
148156
149157def cmd_stream (args : argparse .Namespace ) -> int :
150158 cleaner = StreamingCleaner (** _stream_options (args ))
159+ sanitize = getattr (args , "sanitize_formulas" , False )
151160 return _run_stream (cleaner , _read_chunks (args .input , args .batch_size ),
152- _BatchWriter (args .output ), args .report , args .quiet ,
153- getattr (args , "quarantine" , None ))
161+ _BatchWriter (args .output , sanitize_formulas = sanitize ),
162+ args .report , args .quiet ,
163+ getattr (args , "quarantine" , None ),
164+ sanitize_formulas = sanitize )
154165
155166
156167def cmd_stream_kafka (args : argparse .Namespace ) -> int :
157168 cleaner = StreamingCleaner (** _stream_options (args ))
158169 batches = cleaner .clean_kafka ( # validates the kafka dependency
159170 topic = args .topic , bootstrap_servers = args .bootstrap_servers ,
160171 batch_size = args .batch_size , max_batches = args .max_batches )
161- writer = _BatchWriter (args .output )
172+ writer = _BatchWriter (args .output ,
173+ sanitize_formulas = getattr (args , "sanitize_formulas" , False ))
162174 if args .report :
163175 os .makedirs (args .report , exist_ok = True )
164176 for cleaned , report in batches :
@@ -234,6 +246,9 @@ def add_stream_subparsers(subparsers: argparse._SubParsersAction) -> None:
234246 s = subparsers .add_parser ("stream" , help = "clean a CSV/Parquet file in micro-batches" )
235247 s .add_argument ("input" )
236248 s .add_argument ("-o" , "--output" )
249+ s .add_argument ("--sanitize-formulas" , action = "store_true" ,
250+ help = "prefix ' to string cells starting with = + - @ tab/CR in "
251+ "CSV outputs so spreadsheets render them as text (OWASP)" )
237252 s .add_argument ("--batch-size" , "--chunksize" , type = int , default = 100_000 , dest = "batch_size" )
238253 s .add_argument ("--report" , metavar = "DIR" , help = "directory for per-batch + summary JSON" )
239254 s .add_argument ("--target-column" )
@@ -252,6 +267,9 @@ def add_stream_subparsers(subparsers: argparse._SubParsersAction) -> None:
252267 k .add_argument ("--batch-size" , type = int , default = 10_000 )
253268 k .add_argument ("--max-batches" , type = int )
254269 k .add_argument ("-o" , "--output" )
270+ k .add_argument ("--sanitize-formulas" , action = "store_true" ,
271+ help = "prefix ' to string cells starting with = + - @ tab/CR in "
272+ "CSV outputs so spreadsheets render them as text (OWASP)" )
255273 k .add_argument ("--report" , metavar = "DIR" )
256274 k .add_argument ("--target-column" )
257275 k .add_argument ("--id-columns" , nargs = "*" , default = ())
0 commit comments