Coverage for src/dataknobs_data/pandas/converter.py: 0%

174 statements  

« prev     ^ index     » next       coverage.py v7.10.3, created at 2025-08-17 18:57 -0500

1"""Core converter between DataKnobs Records and Pandas DataFrames.""" 

2 

3from dataclasses import dataclass 

4from typing import Any, Dict, List, Optional, Union 

5 

6import pandas as pd 

7import numpy as np 

8 

9from dataknobs_data.records import Record 

10from dataknobs_data.fields import Field, FieldType 

11from .type_mapper import TypeMapper 

12from .metadata import MetadataHandler, MetadataConfig, MetadataStrategy 

13 

14 

15@dataclass 

16class ConversionOptions: 

17 """Options for conversion between Records and DataFrames.""" 

18 include_metadata: bool = False 

19 metadata_columns: List[str] = None # Columns to treat as metadata 

20 flatten_nested: bool = False # Flatten nested structures 

21 preserve_index: bool = True 

22 use_index_as_id: bool = False # Use DataFrame index as record ID 

23 type_mapping: Dict[str, str] = None # Custom type mappings 

24 null_handling: str = "preserve" # "preserve", "drop", "fill" 

25 datetime_format: Optional[str] = None # Format for datetime conversion 

26 timezone: Optional[str] = None # Timezone for datetime conversion 

27 

28 # Keep these for backward compatibility 

29 preserve_types: bool = True 

30 index_column: Optional[str] = None # Use specific field as index 

31 flatten_json: bool = False 

32 metadata_strategy: MetadataStrategy = MetadataStrategy.ATTRS 

33 handle_missing: str = "preserve" # "preserve", "drop", "fill" 

34 fill_value: Any = None 

35 

36 def __post_init__(self): 

37 """Initialize default values for mutable parameters.""" 

38 if self.metadata_columns is None: 

39 self.metadata_columns = [] 

40 if self.type_mapping is None: 

41 self.type_mapping = {} 

42 

43 def merge_metadata(self, meta1: Dict[str, Any], meta2: Dict[str, Any]) -> Dict[str, Any]: 

44 """Merge two metadata dictionaries. 

45  

46 Args: 

47 meta1: First metadata dict 

48 meta2: Second metadata dict (overwrites meta1 on conflicts) 

49  

50 Returns: 

51 Merged metadata dictionary 

52 """ 

53 result = meta1.copy() 

54 for key, value in meta2.items(): 

55 if key in result and isinstance(result[key], dict) and isinstance(value, dict): 

56 # Recursively merge nested dicts 

57 result[key] = self.merge_metadata(result[key], value) 

58 else: 

59 result[key] = value 

60 return result 

61 

62 

63class DataFrameConverter: 

64 """Converts between DataKnobs Records and Pandas DataFrames.""" 

65 

66 def __init__(self, type_mapper: Optional[TypeMapper] = None): 

67 """Initialize converter. 

68  

69 Args: 

70 type_mapper: Custom type mapper (uses default if None) 

71 """ 

72 self.type_mapper = type_mapper or TypeMapper() 

73 

74 def records_to_dataframe( 

75 self, 

76 records: List[Record], 

77 options: Optional[ConversionOptions] = None 

78 ) -> pd.DataFrame: 

79 """Convert list of Records to DataFrame. 

80  

81 Args: 

82 records: List of Records to convert 

83 options: Conversion options 

84  

85 Returns: 

86 Pandas DataFrame 

87 """ 

88 options = options or ConversionOptions() 

89 

90 if not records: 

91 return pd.DataFrame() 

92 

93 # Extract data from records 

94 data_rows = [] 

95 for record in records: 

96 row = {} 

97 

98 # Add field values 

99 for field_name, field in record.fields.items(): 

100 if options.flatten_nested and isinstance(field.value, dict): 

101 # Flatten nested dictionaries 

102 for nested_key, nested_val in field.value.items(): 

103 row[f"{field_name}.{nested_key}"] = nested_val 

104 else: 

105 row[field_name] = field.value 

106 

107 # Add metadata as a column if requested 

108 if options.include_metadata and record.metadata: 

109 row["_metadata"] = record.metadata 

110 

111 data_rows.append(row) 

112 

113 # Create DataFrame 

114 df = pd.DataFrame(data_rows) 

115 

116 # Preserve column order from the first record for consistency 

117 # This maintains the order fields were added to records 

118 if not df.empty and records: 

119 # Get column order from first record's fields 

120 first_record = records[0] 

121 column_order = [] 

122 

123 # Add field columns in the order they appear in the record 

124 for field_name in first_record.fields.keys(): 

125 if options.flatten_nested and isinstance(first_record.fields[field_name].value, dict): 

126 # Add flattened columns 

127 for nested_key in first_record.fields[field_name].value.keys(): 

128 col_name = f"{field_name}.{nested_key}" 

129 if col_name in df.columns: 

130 column_order.append(col_name) 

131 elif field_name in df.columns: 

132 column_order.append(field_name) 

133 

134 # Add any remaining columns (like _metadata) at the end 

135 for col in df.columns: 

136 if col not in column_order: 

137 column_order.append(col) 

138 

139 # Reorder DataFrame columns 

140 df = df[column_order] 

141 

142 # Set index if specified 

143 if options.index_column and options.index_column in df.columns: 

144 df = df.set_index(options.index_column) 

145 elif options.preserve_index: 

146 # Only use record IDs as index if: 

147 # 1. They exist 

148 # 2. They're not coming from a data field (id or record_id columns) 

149 # 3. They don't look like auto-generated UUIDs 

150 record_ids = [r.id for r in records] 

151 

152 # Check if the IDs are from data fields 

153 ids_from_fields = any( 

154 'id' in r.fields or 'record_id' in r.fields 

155 for r in records 

156 ) 

157 

158 # Only set index if IDs exist, aren't from fields, and aren't UUIDs 

159 if (any(record_ids) and not ids_from_fields and not all( 

160 id and len(id) == 36 and id.count('-') == 4 

161 for id in record_ids if id 

162 )): 

163 df.index = record_ids 

164 df.index.name = "record_id" 

165 

166 return df 

167 

168 def dataframe_to_records( 

169 self, 

170 df: pd.DataFrame, 

171 options: Optional[ConversionOptions] = None 

172 ) -> List[Record]: 

173 """Convert DataFrame to list of Records. 

174  

175 Args: 

176 df: DataFrame to convert 

177 options: Conversion options 

178  

179 Returns: 

180 List of Records 

181 """ 

182 options = options or ConversionOptions() 

183 

184 records = [] 

185 

186 # Convert each row to a Record 

187 for idx, row in df.iterrows(): 

188 # Extract metadata for this row from metadata columns 

189 row_metadata = {} 

190 if options.metadata_columns: 

191 for col in options.metadata_columns: 

192 if col in row.index: 

193 col_value = row[col] 

194 # If the column value is a dict, merge it into metadata 

195 if isinstance(col_value, dict): 

196 row_metadata.update(col_value) 

197 else: 

198 # Otherwise store it with the column name (without leading underscore) 

199 row_metadata[col.lstrip('_')] = col_value 

200 

201 # Prepare row data (excluding metadata columns) 

202 row_data = {} 

203 for col in row.index: 

204 if col not in options.metadata_columns: 

205 row_data[col] = row[col] 

206 

207 # Determine record ID 

208 record_id = None 

209 if options.use_index_as_id: 

210 if isinstance(idx, str): 

211 record_id = idx 

212 elif idx is not None and not pd.isna(idx): 

213 record_id = str(idx) 

214 

215 # Create record 

216 record = Record(data=row_data, metadata=row_metadata, id=record_id) 

217 records.append(record) 

218 

219 return records 

220 

221 def record_to_series(self, record: Record) -> pd.Series: 

222 """Convert a single Record to a Pandas Series. 

223  

224 Args: 

225 record: Record to convert 

226  

227 Returns: 

228 Pandas Series 

229 """ 

230 data = {} 

231 for field_name, field in record.fields.items(): 

232 value = self.type_mapper.convert_value_to_pandas(field.value, field.type) 

233 data[field_name] = value 

234 

235 series = pd.Series(data) 

236 if record.id: 

237 series.name = record.id 

238 

239 return series 

240 

241 def series_to_record( 

242 self, 

243 series: pd.Series, 

244 record_id: Optional[str] = None 

245 ) -> Record: 

246 """Convert a Pandas Series to a Record. 

247  

248 Args: 

249 series: Series to convert 

250 record_id: Optional record ID 

251  

252 Returns: 

253 Record 

254 """ 

255 record = Record(id=record_id or series.name if hasattr(series, 'name') else None) 

256 

257 for column, value in series.items(): 

258 # Skip metadata columns 

259 if isinstance(column, str) and column.startswith("_meta_"): 

260 continue 

261 

262 # Infer field type 

263 field_type = self.type_mapper.infer_field_type_from_value(value) 

264 

265 # Convert value 

266 field_value = self.type_mapper.convert_value_from_pandas(value, field_type) 

267 

268 # Create field 

269 field = Field( 

270 name=str(column), 

271 value=field_value, 

272 type=field_type 

273 ) 

274 record.fields[str(column)] = field 

275 

276 return record 

277 

278 def _series_to_record( 

279 self, 

280 row: pd.Series, 

281 idx: Any, 

282 options: ConversionOptions 

283 ) -> Record: 

284 """Convert a DataFrame row to a Record. 

285  

286 Args: 

287 row: DataFrame row as Series 

288 idx: Row index 

289 options: Conversion options 

290  

291 Returns: 

292 Record 

293 """ 

294 # Determine record ID 

295 record_id = None 

296 if options.preserve_index: 

297 if isinstance(idx, str): 

298 record_id = idx 

299 elif idx is not None and not pd.isna(idx): 

300 record_id = str(idx) 

301 

302 record = Record(id=record_id) 

303 

304 for column, value in row.items(): 

305 # Skip metadata columns 

306 if isinstance(column, str) and column.startswith("_meta_"): 

307 continue 

308 

309 # Handle multi-index columns 

310 if isinstance(column, tuple): 

311 column = column[0] # Use first level 

312 

313 # Infer field type 

314 field_type = self.type_mapper.infer_field_type_from_value(value) 

315 

316 # Convert value 

317 field_value = self.type_mapper.convert_value_from_pandas(value, field_type) 

318 

319 # Create field 

320 field = Field( 

321 name=str(column), 

322 value=field_value, 

323 type=field_type 

324 ) 

325 record.fields[str(column)] = field 

326 

327 return record 

328 

329 def _flatten_json_value(self, value: Any) -> Any: 

330 """Flatten JSON value for DataFrame insertion. 

331  

332 Args: 

333 value: JSON value (dict or list) 

334  

335 Returns: 

336 Flattened value or string representation 

337 """ 

338 # Check for None explicitly first 

339 if value is None: 

340 return value 

341 

342 # Check for pandas NA types 

343 try: 

344 if pd.isna(value): 

345 return value 

346 except (TypeError, ValueError): 

347 # pd.isna doesn't work with lists/dicts 

348 pass 

349 

350 if isinstance(value, dict): 

351 # For dict, could expand to multiple columns 

352 # For now, convert to string 

353 return str(value) 

354 elif isinstance(value, list): 

355 # For list, convert to string 

356 return str(value) 

357 

358 return value 

359 

360 def validate_conversion( 

361 self, 

362 records: List[Record], 

363 df: pd.DataFrame, 

364 options: Optional[ConversionOptions] = None 

365 ) -> Dict[str, Any]: 

366 """Validate conversion accuracy. 

367  

368 Args: 

369 records: Original records 

370 df: Converted DataFrame 

371 options: Conversion options used 

372  

373 Returns: 

374 Validation report 

375 """ 

376 options = options or ConversionOptions() 

377 

378 report = { 

379 "record_count_match": len(records) == len(df), 

380 "original_record_count": len(records), 

381 "dataframe_row_count": len(df), 

382 "field_preservation": {}, 

383 "type_preservation": {}, 

384 "value_accuracy": {} 

385 } 

386 

387 # Check field preservation 

388 original_fields = set() 

389 for record in records: 

390 original_fields.update(record.fields.keys()) 

391 

392 df_columns = set(df.columns) 

393 if options.metadata_strategy == MetadataStrategy.COLUMNS: 

394 df_columns = {col for col in df_columns if not col.startswith("_meta_")} 

395 

396 report["field_preservation"] = { 

397 "original_fields": sorted(original_fields), 

398 "dataframe_columns": sorted(df_columns), 

399 "missing_fields": sorted(original_fields - df_columns), 

400 "extra_columns": sorted(df_columns - original_fields) 

401 } 

402 

403 # Check type preservation if enabled 

404 if options.preserve_types: 

405 for record in records[:10]: # Sample first 10 records 

406 for field_name, field in record.fields.items(): 

407 if field_name in df.columns: 

408 df_dtype = str(df[field_name].dtype) 

409 expected_dtype = str(self.type_mapper.field_type_to_pandas(field.type)) 

410 if df_dtype != expected_dtype: 

411 report["type_preservation"][field_name] = { 

412 "expected": expected_dtype, 

413 "actual": df_dtype 

414 } 

415 

416 return report