@@ -65,28 +65,11 @@ def __init__(
6565 auth_token : str | None = None ,
6666 ):
6767 self ._gateway = None
68+ self ._is_gateway_version_checked = False
6869 self .address = address or configuration .JAVA_GATEWAY_ADDRESS
6970 self .port = port or configuration .JAVA_GATEWAY_PORT
7071 self .auto_convert = auto_convert or configuration .JAVA_GATEWAY_AUTO_CONVERT
7172 self .auth_token = auth_token or configuration .JAVA_GATEWAY_AUTH_TOKEN
72- gateway_version = "unknown"
73- with contextlib .suppress (Py4JError ):
74- # 1. Java gateway version is too old: doesn't have method 'getGatewayVersion()'
75- # 2. Error connecting to Java gateway
76- gateway_version = self .get_gateway_version ()
77- if (
78- not __version__ .endswith ("dev" )
79- and gateway_version
80- and not version_match (Version .DS , gateway_version )
81- ):
82- warnings .warn (
83- f"Using unmatched version of pydolphinscheduler (version { __version__ } ) "
84- f"and Java gateway (version { gateway_version } ) may cause errors. "
85- "We strongly recommend you to find the matched version "
86- "(check: https://pypi.org/project/apache-dolphinscheduler)" ,
87- UserWarning ,
88- stacklevel = 2 ,
89- )
9073
9174 @property
9275 def gateway (self ) -> JavaGateway :
@@ -108,6 +91,35 @@ def gateway(self) -> JavaGateway:
10891 self ._gateway = JavaGateway (gateway_parameters = gateway_parameters )
10992 return self ._gateway
11093
94+ def _check_gateway_version (self ):
95+ """Warn once when Python SDK and Java gateway versions do not match."""
96+ if self ._is_gateway_version_checked :
97+ return
98+ self ._is_gateway_version_checked = True
99+ if __version__ .endswith ("dev" ):
100+ return
101+
102+ gateway_version = "unknown"
103+ with contextlib .suppress (Py4JError ):
104+ # 1. Java gateway version is too old: doesn't have method 'getGatewayVersion()'
105+ # 2. Error connecting to Java gateway
106+ gateway_version = self .gateway .entry_point .getGatewayVersion ()
107+ if gateway_version and not version_match (Version .DS , gateway_version ):
108+ warnings .warn (
109+ f"Using unmatched version of pydolphinscheduler (version { __version__ } ) "
110+ f"and Java gateway (version { gateway_version } ) may cause errors. "
111+ "We strongly recommend you to find the matched version "
112+ "(check: https://pypi.org/project/apache-dolphinscheduler)" ,
113+ UserWarning ,
114+ stacklevel = 2 ,
115+ )
116+
117+ @property
118+ def entry_point (self ):
119+ """Return Java gateway entry point with lazy version validation."""
120+ self ._check_gateway_version ()
121+ return self .gateway .entry_point
122+
111123 def get_gateway_version (self ):
112124 """Get the java gateway version, expected to be equal with pydolphinscheduler."""
113125 return self .gateway .entry_point .getGatewayVersion ()
@@ -120,69 +132,69 @@ def get_datasource(self, name: str, type: str | None = None):
120132 :param name: datasource name of the datasource to be queried
121133 :param type: datasource type of the datasource, only used to filter the result.
122134 """
123- return self .gateway . entry_point .getDatasource (name , type )
135+ return self .entry_point .getDatasource (name , type )
124136
125137 def get_resources_file_info (self , program_type : str , main_package : str ):
126138 """Get resources file info through java gateway."""
127- return self .gateway . entry_point .getResourcesFileInfo (program_type , main_package )
139+ return self .entry_point .getResourcesFileInfo (program_type , main_package )
128140
129141 def create_or_update_resource (self , user_name : str , name : str , content : str ):
130142 """Create or update resource through java gateway."""
131- return self .gateway . entry_point .createOrUpdateResource (user_name , name , content )
143+ return self .entry_point .createOrUpdateResource (user_name , name , content )
132144
133145 def query_resources_file_info (self , user_name : str , name : str ):
134146 """Get resources file info through java gateway."""
135- return self .gateway . entry_point .queryResourcesFileInfo (user_name , name )
147+ return self .entry_point .queryResourcesFileInfo (user_name , name )
136148
137149 def query_environment_info (self , name : str ):
138150 """Get environment info through java gateway."""
139- return self .gateway . entry_point .getEnvironmentInfo (name )
151+ return self .entry_point .getEnvironmentInfo (name )
140152
141153 def get_code_and_version (
142154 self , project_name : str , workflow_name : str , task_name : str
143155 ):
144156 """Get code and version through java gateway."""
145- return self .gateway . entry_point .getCodeAndVersion (
157+ return self .entry_point .getCodeAndVersion (
146158 project_name , workflow_name , task_name
147159 )
148160
149161 def create_or_grant_project (
150162 self , user : str , name : str , description : str | None = None
151163 ):
152164 """Create or grant project through java gateway."""
153- return self .gateway . entry_point .createOrGrantProject (user , name , description )
165+ return self .entry_point .createOrGrantProject (user , name , description )
154166
155167 def query_project_by_name (self , user : str , name : str ):
156168 """Query project through java gateway."""
157- return self .gateway . entry_point .queryProjectByName (user , name )
169+ return self .entry_point .queryProjectByName (user , name )
158170
159171 def update_project (
160172 self , user : str , project_code : int , project_name : str , description : str
161173 ):
162174 """Update project through java gateway."""
163- return self .gateway . entry_point .updateProject (
175+ return self .entry_point .updateProject (
164176 user , project_code , project_name , description
165177 )
166178
167179 def delete_project (self , user : str , code : int ):
168180 """Delete project through java gateway."""
169- return self .gateway . entry_point .deleteProject (user , code )
181+ return self .entry_point .deleteProject (user , code )
170182
171183 def create_tenant (
172184 self , tenant_name : str , queue_name : str , description : str | None = None
173185 ):
174186 """Create tenant through java gateway."""
175- return self .gateway . entry_point .createTenant (
187+ return self .entry_point .createTenant (
176188 tenant_name , description , queue_name
177189 )
178190
179191 def query_tenant (self , tenant_code : str ):
180192 """Query tenant through java gateway."""
181- return self .gateway . entry_point .queryTenantByCode (tenant_code )
193+ return self .entry_point .queryTenantByCode (tenant_code )
182194
183195 def grant_tenant_to_user (self , user_name : str , tenant_code : str ):
184196 """Grant tenant to user through java gateway."""
185- return self .gateway . entry_point .grantTenantToUser (user_name , tenant_code )
197+ return self .entry_point .grantTenantToUser (user_name , tenant_code )
186198
187199 def update_tenant (
188200 self ,
@@ -193,13 +205,13 @@ def update_tenant(
193205 description : str | None = None ,
194206 ):
195207 """Update tenant through java gateway."""
196- return self .gateway . entry_point .updateTenant (
208+ return self .entry_point .updateTenant (
197209 user , tenant_id , code , queue_id , description
198210 )
199211
200212 def delete_tenant (self , user : str , tenant_id : int ):
201213 """Delete tenant through java gateway."""
202- return self .gateway . entry_point .deleteTenantById (user , tenant_id )
214+ return self .entry_point .deleteTenantById (user , tenant_id )
203215
204216 def create_user (
205217 self ,
@@ -212,7 +224,7 @@ def create_user(
212224 status : int ,
213225 ):
214226 """Create user through java gateway."""
215- return self .gateway . entry_point .createUser (
227+ return self .entry_point .createUser (
216228 name , password , email , phone , tenant , queue , status
217229 )
218230
@@ -231,13 +243,13 @@ def update_user(
231243 status : int ,
232244 ):
233245 """Update user through java gateway."""
234- return self .gateway . entry_point .updateUser (
246+ return self .entry_point .updateUser (
235247 name , password , email , phone , tenant , queue , status
236248 )
237249
238250 def delete_user (self , name : str , user_id : int ):
239251 """Delete user through java gateway."""
240- return self .gateway . entry_point .deleteUser (name , user_id )
252+ return self .entry_point .deleteUser (name , user_id )
241253
242254 def get_dependent_info (
243255 self ,
@@ -246,13 +258,13 @@ def get_dependent_info(
246258 task_name : str | None = None ,
247259 ):
248260 """Get dependent info through java gateway."""
249- return self .gateway . entry_point .getDependentInfo (
261+ return self .entry_point .getDependentInfo (
250262 project_name , workflow_name , task_name
251263 )
252264
253265 def get_workflow_info (self , user_name : str , project_name : str , workflow_name : str ):
254266 """Get workflow info through java gateway."""
255- return self .gateway . entry_point .getWorkflowInfo (
267+ return self .entry_point .getWorkflowInfo (
256268 user_name , project_name , workflow_name
257269 )
258270
@@ -276,7 +288,7 @@ def create_or_update_workflow(
276288 other_params_json : str | None = None ,
277289 ):
278290 """Create or update workflow through java gateway."""
279- return self .gateway . entry_point .createOrUpdateWorkflow (
291+ return self .entry_point .createOrUpdateWorkflow (
280292 user_name ,
281293 project_name ,
282294 name ,
@@ -305,7 +317,7 @@ def exec_workflow_instance(
305317 warning_group_id : int ,
306318 ):
307319 """Exec workflow instance through java gateway."""
308- return self .gateway . entry_point .execWorkflowInstance (
320+ return self .entry_point .execWorkflowInstance (
309321 user_name ,
310322 project_name ,
311323 workflow_name ,
0 commit comments