Test RayDedupNode functionality.
(
ray_session, ray_config, mock_context, sample_data, column_info
)
| 277 | |
| 278 | |
| 279 | def test_ray_dedup_node( |
| 280 | ray_session, ray_config, mock_context, sample_data, column_info |
| 281 | ): |
| 282 | """Test RayDedupNode functionality.""" |
| 283 | duplicated_data = pd.concat([sample_data, sample_data.iloc[:1]], ignore_index=True) |
| 284 | ray_dataset = ray.data.from_pandas(duplicated_data) |
| 285 | input_value = DAGValue(data=ray_dataset, format=DAGFormat.RAY) |
| 286 | dummy_node = DummyInputNode("input_node", input_value) |
| 287 | node = RayDedupNode( |
| 288 | name="dedup", |
| 289 | column_info=column_info, |
| 290 | config=ray_config, |
| 291 | ) |
| 292 | node.add_input(dummy_node) |
| 293 | mock_context.node_outputs = {"input_node": input_value} |
| 294 | result = node.execute(mock_context) |
| 295 | assert isinstance(result, DAGValue) |
| 296 | assert result.format == DAGFormat.RAY |
| 297 | result_df = result.data.to_pandas() |
| 298 | assert len(result_df) == 2 # Should remove the duplicate row |
| 299 | assert "driver_id" in result_df.columns |
| 300 | |
| 301 | |
| 302 | def test_ray_dedup_node_materialization_within_block( |
nothing calls this directly
no test coverage detected