MCPcopy Create free account
hub / github.com/feast-dev/feast / RayReadNode

Class RayReadNode

sdk/python/feast/infra/compute_engines/ray/nodes.py:40–101  ·  view source on GitHub ↗

Ray node for reading data from offline stores.

Source from the content-addressed store, hash-verified

38
39
40class RayReadNode(DAGNode):
41 """
42 Ray node for reading data from offline stores.
43 """
44
45 def __init__(
46 self,
47 name: str,
48 source: DataSource,
49 column_info,
50 config: RayComputeEngineConfig,
51 start_time: Optional[datetime] = None,
52 end_time: Optional[datetime] = None,
53 ):
54 super().__init__(name)
55 self.source = source
56 self.column_info = column_info
57 self.config = config
58 self.start_time = start_time
59 self.end_time = end_time
60
61 def execute(self, context: ExecutionContext) -> DAGValue:
62 """Execute the read operation to load data from the offline store."""
63 try:
64 retrieval_job = create_offline_store_retrieval_job(
65 data_source=self.source,
66 column_info=self.column_info,
67 context=context,
68 start_time=self.start_time,
69 end_time=self.end_time,
70 )
71
72 if hasattr(retrieval_job, "to_ray_dataset"):
73 ray_dataset = retrieval_job.to_ray_dataset()
74 else:
75 try:
76 arrow_table = retrieval_job.to_arrow()
77 ray_wrapper = get_ray_wrapper()
78 ray_dataset = ray_wrapper.from_arrow(arrow_table)
79 except Exception:
80 df = retrieval_job.to_df()
81 ray_wrapper = get_ray_wrapper()
82 ray_dataset = ray_wrapper.from_pandas(df)
83
84 field_mapping = getattr(self.source, "field_mapping", None)
85 if field_mapping:
86 ray_dataset = apply_field_mapping(ray_dataset, field_mapping)
87
88 return DAGValue(
89 data=ray_dataset,
90 format=DAGFormat.RAY,
91 metadata={
92 "source": "offline_store",
93 "source_type": type(self.source).__name__,
94 "start_time": self.start_time,
95 "end_time": self.end_time,
96 },
97 )

Callers 2

build_source_nodeMethod · 0.90
test_ray_read_nodeFunction · 0.90

Calls

no outgoing calls

Tested by 1

test_ray_read_nodeFunction · 0.72