FazBrowse GitHub Viewer
|
Trending
|
URL:
|
Home
Tools:
[Download Repo ZIP]
[View Raw Code]
[Original HTTPS Page]
feast/sdk/python/feast/arrow_error_handler.py at v0.63.0 · feast-dev/feast · GitHub
Uh oh!
There was an error while loading.
Please reload this page
.
feast-dev
/
feast
Public
Notifications
You must be signed in to change notification settings
Fork
1.4k
Star
7.2k
Code
Issues
215
Pull requests
190
Discussions
Actions
Security and quality
1
Insights
Additional navigation options
Code
Issues
Pull requests
Discussions
Actions
Security and quality
Insights
Expand file tree
Breadcrumbs
feast
/
sdk
/
python
/
feast
/
arrow_error_handler.py
Copy path
More file actions
More file actions
Latest commit
History
History
History
94 lines (78 loc) · 2.97 KB
Breadcrumbs
feast
/
sdk
/
python
/
feast
/
arrow_error_handler.py
Copy path
File metadata and controls
94 lines (78 loc) · 2.97 KB
Raw
Copy raw file
Download raw file
Open symbols panel
Edit and raw actions
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
import
json
import
logging
import
time
from
functools
import
wraps
import
pyarrow
.
flight
as
fl
from
feast
.
errors
import
FeastError
logger
=
logging
.
getLogger
(
__name__
)
BACKOFF_FACTOR
=
0.5
def
arrow_client_error_handling_decorator
(
func
):
@
wraps
(
func
)
def
wrapper
(
*
args
,
**
kwargs
):
# Retry only applies to FeastFlightClient methods where args[0] (self)
# carries _connection_retries from RemoteOfflineStoreConfig.
# Standalone stream functions (write_table, read_all) get 0 retries:
# broken streams can't be reused and retrying risks duplicate writes.
max_retries
=
max
(
0
,
getattr
(
args
[
0
],
"_connection_retries"
,
0
))
if
args
else
0
for
attempt
in
range
(
max_retries
+
1
):
try
:
return
func
(
*
args
,
**
kwargs
)
except
fl
.
FlightUnavailableError
as
e
:
if
attempt
<
max_retries
:
wait_time
=
BACKOFF_FACTOR
*
(
2
**
attempt
)
logger
.
warning
(
"Transient Arrow Flight error on attempt %d/%d, "
"retrying in %.1fs: %s"
,
attempt
+
1
,
max_retries
+
1
,
wait_time
,
e
,
)
time
.
sleep
(
wait_time
)
continue
mapped_error
=
FeastError
.
from_error_detail
(
_get_exception_data
(
str
(
e
)))
if
mapped_error
is
not
None
:
raise
mapped_error
raise
e
except
Exception
as
e
:
mapped_error
=
FeastError
.
from_error_detail
(
_get_exception_data
(
e
.
args
[
0
])
)
if
mapped_error
is
not
None
:
raise
mapped_error
raise
e
return
wrapper
def
arrow_server_error_handling_decorator
(
func
):
@
wraps
(
func
)
def
wrapper
(
*
args
,
**
kwargs
):
try
:
return
func
(
*
args
,
**
kwargs
)
except
Exception
as
e
:
if
isinstance
(
e
,
FeastError
):
raise
fl
.
FlightError
(
e
.
to_error_detail
())
# Re-raise non-Feast exceptions so Arrow Flight returns a proper error
# instead of allowing the server method to return None.
raise
e
return
wrapper
def
_get_exception_data
(
except_str
)
->
str
:
if
not
isinstance
(
except_str
,
str
):
return
""
substring
=
"Flight error: "
position
=
except_str
.
find
(
substring
)
if
position
==
-
1
:
return
""
search_start
=
position
+
len
(
substring
)
json_start
=
except_str
.
find
(
"{"
,
search_start
)
if
json_start
==
-
1
:
return
""
search_pos
=
json_start
while
True
:
json_end
=
except_str
.
find
(
"}"
,
search_pos
+
1
)
if
json_end
==
-
1
:
return
""
candidate
=
except_str
[
json_start
:
json_end
+
1
]
try
:
json
.
loads
(
candidate
)
return
candidate
except
(
json
.
JSONDecodeError
,
ValueError
):
search_pos
=
json_end
Back
|
FazBrowse Home
|
New Git URL