Test RayReadNode functionality.
(ray_session, ray_config, mock_context, sample_data, column_info)
| 141 | |
| 142 | |
| 143 | def test_ray_read_node(ray_session, ray_config, mock_context, sample_data, column_info): |
| 144 | """Test RayReadNode functionality.""" |
| 145 | ray_dataset = ray.data.from_pandas(sample_data) |
| 146 | mock_source = DummySource() |
| 147 | node = RayReadNode( |
| 148 | name="read", |
| 149 | source=mock_source, |
| 150 | column_info=column_info, |
| 151 | config=ray_config, |
| 152 | ) |
| 153 | mock_context.registry = None |
| 154 | mock_context.store = None |
| 155 | mock_context.offline_store = None |
| 156 | mock_retrieval_job = DummyRetrievalJob(ray_dataset) |
| 157 | import feast.infra.compute_engines.ray.nodes as ray_nodes |
| 158 | |
| 159 | ray_nodes.create_offline_store_retrieval_job = lambda **kwargs: mock_retrieval_job |
| 160 | result = node.execute(mock_context) |
| 161 | assert isinstance(result, DAGValue) |
| 162 | assert result.format == DAGFormat.RAY |
| 163 | result_df = result.data.to_pandas() |
| 164 | assert len(result_df) == 3 |
| 165 | assert "driver_id" in result_df.columns |
| 166 | assert "conv_rate" in result_df.columns |
| 167 | |
| 168 | |
| 169 | def test_ray_aggregation_node( |
nothing calls this directly
no test coverage detected