Skip to content

Commit 72cfc48

Browse files
committed
refactor import/sync drop dtable-server view-rows
1 parent 5260b27 commit 72cfc48

7 files changed

Lines changed: 757 additions & 423 deletions

File tree

dtable_events/common_dataset/common_dataset_sync_utils.py

Lines changed: 329 additions & 247 deletions
Large diffs are not rendered by default.

dtable_events/common_dataset/common_dataset_syncer.py

Lines changed: 60 additions & 99 deletions
Original file line numberDiff line numberDiff line change
@@ -9,10 +9,9 @@
99

1010
from dtable_events import init_db_session_class
1111
from dtable_events.app.config import DTABLE_PRIVATE_KEY
12-
from dtable_events.common_dataset.common_dataset_sync_utils import import_or_sync, set_common_dataset_invalid, set_common_dataset_sync_invalid
12+
from dtable_events.common_dataset.common_dataset_sync_utils import import_sync_CDS, set_common_dataset_invalid, set_common_dataset_sync_invalid
1313
from dtable_events.utils import get_opt_from_conf_or_env, parse_bool, uuid_str_to_36_chars, get_inner_dtable_server_url
14-
15-
logger = logging.getLogger(__name__)
14+
from dtable_events.utils.dtable_server_api import DTableServerAPI
1615

1716
class CommonDatasetSyncer(object):
1817

@@ -54,7 +53,7 @@ def get_dtable_server_header(dtable_uuid):
5453
algorithm='HS256'
5554
)
5655
except Exception as e:
57-
logger.error(e)
56+
logging.error(e)
5857
return
5958
return {'Authorization': 'Token ' + access_token}
6059

@@ -63,89 +62,62 @@ def gen_src_dst_assets(dst_dtable_uuid, src_dtable_uuid, src_table_id, src_view_
6362
"""
6463
return assets -> dict
6564
"""
66-
dst_headers = get_dtable_server_header(dst_dtable_uuid)
67-
src_headers = get_dtable_server_header(src_dtable_uuid)
68-
69-
# request src_dtable
7065
dtable_server_url = get_inner_dtable_server_url()
71-
url = dtable_server_url.strip('/') + '/dtables/' + src_dtable_uuid + '?from=dtable_events'
72-
66+
src_dtable_server_api = DTableServerAPI('dtable-events', src_dtable_uuid, dtable_server_url)
67+
dst_dtable_server_api = DTableServerAPI('dtable-events', dst_dtable_uuid, dtable_server_url)
7368
try:
74-
resp = requests.get(url, headers=src_headers)
75-
src_dtable_json = resp.json()
69+
src_dtable_metadata = src_dtable_server_api.get_metadata()
70+
dst_dtable_metadata = dst_dtable_server_api.get_metadata()
7671
except Exception as e:
77-
logger.error('request src dtable: %s error: %s', src_dtable_uuid, e)
78-
return
72+
logging.error('request src dst dtable: %s, %s metadata error: %s', src_dtable_uuid, dst_dtable_uuid, e)
73+
return None
7974

80-
# check src_table src_view
81-
src_table = None
82-
for table in src_dtable_json.get('tables', []):
83-
if table.get('_id') == src_table_id:
75+
src_table, src_view = None, None
76+
for table in src_dtable_metadata.get('tables', []):
77+
if table['_id'] == src_table_id:
8478
src_table = table
8579
break
8680
if not src_table:
8781
set_common_dataset_invalid(dataset_id, db_session)
8882
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
8983
logging.error('Source table not found.')
90-
return
91-
92-
src_view = None
93-
if src_view_id:
94-
for view in src_table.get('views', []):
95-
if view.get('_id') == src_view_id:
96-
src_view = view
97-
break
98-
if not src_view:
99-
set_common_dataset_invalid(dataset_id, db_session)
100-
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
101-
logging.error('Source view not found.')
102-
return
103-
else:
104-
views = src_table.get('views', [])
105-
if not views or not isinstance(views, list):
106-
set_common_dataset_invalid(dataset_id, db_session)
107-
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
108-
logging.error('No views found.')
109-
return
110-
src_view = views[0]
84+
return None
85+
for view in src_table.get('views', []):
86+
if view['_id'] == src_view_id:
87+
src_view = view
88+
break
89+
if not src_view:
90+
set_common_dataset_invalid(dataset_id, db_session)
91+
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
92+
logging.error('Source view not found.')
93+
return None
11194

112-
# get src columns
113-
src_view_hidden_columns = src_view.get('hidden_columns', [])
114-
if not src_view_hidden_columns:
115-
src_columns = src_table.get('columns', [])
116-
else:
117-
src_columns = [col for col in src_table.get('columns', []) if col.get('key') not in src_view_hidden_columns]
95+
src_columns = [col for col in src_table.get('columns', []) if col not in src_view.get('hidden_columns', [])]
11896

119-
# request dst_dtable
120-
url = dtable_server_url.strip('/') + '/dtables/' + dst_dtable_uuid + '?from=dtable_events'
121-
try:
122-
resp = requests.get(url, headers=dst_headers)
123-
dst_dtable_json = resp.json()
124-
except Exception as e:
125-
logging.error('request dst dtable: %s error: %s', dst_dtable_uuid, e)
126-
return
97+
src_enable_archive = (src_dtable_metadata.get('settings') or {}).get('enable_archive', False)
98+
src_version = src_dtable_metadata.get('version')
12799

128-
# check dst_table
129100
dst_table = None
130-
for table in dst_dtable_json.get('tables', []):
131-
if table.get('_id') == dst_table_id:
132-
dst_table = table
133-
break
134-
if not dst_table:
135-
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
136-
logging.warning('Destination table: %s not found.' % dst_table_id)
137-
return
101+
if dst_table_id:
102+
for table in dst_dtable_metadata.get('tables', []):
103+
if table['_id'] == dst_table_id:
104+
dst_table = table
105+
break
106+
if not dst_table:
107+
set_common_dataset_invalid(dataset_id, db_session)
108+
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
109+
logging.error('Destination table not found.')
110+
return None
138111

139112
return {
140-
'dst_headers': dst_headers,
141-
'src_headers': src_headers,
142-
'src_table': src_table,
143-
'src_view': src_view,
113+
'src_table_name': src_table['name'],
114+
'src_view_name': src_view['name'],
115+
'src_view_type': src_view.get('type', 'table'),
144116
'src_columns': src_columns,
145-
'dst_columns': dst_table.get('columns'),
146-
'dst_rows': dst_table.get('rows'),
147-
'dst_table_name': dst_table.get('name'),
148-
'src_version': src_dtable_json.get('version')
117+
'src_enable_archive': src_enable_archive,
118+
'src_version': src_version,
119+
'dst_table_name': dst_table['name'] if dst_table else None,
120+
'dst_columns': dst_table['columns'] if dst_table else None
149121
}
150122

151123

@@ -207,61 +179,50 @@ def check_common_dataset(db_session):
207179
if not assets:
208180
continue
209181

210-
dst_headers = assets.get('dst_headers')
211-
src_table = assets.get('src_table')
212-
src_view = assets.get('src_view')
213-
src_columns = assets.get('src_columns')
214-
src_headers = assets.get('src_headers')
215-
dst_columns = assets.get('dst_columns')
216-
dst_rows = assets.get('dst_rows')
217-
dst_table_name = assets.get('dst_table_name')
218-
dtable_src_version = assets.get('src_version')
219-
220-
if dtable_src_version == last_src_version:
182+
if assets.get('src_version') == last_src_version:
221183
continue
222184

223185
try:
224-
result = import_or_sync({
225-
'dst_dtable_uuid': dst_dtable_uuid,
186+
result = import_sync_CDS({
226187
'src_dtable_uuid': src_dtable_uuid,
227-
'src_rows': src_table.get('rows', []),
228-
'src_columns': src_columns,
229-
'src_table_name': src_table.get('name'),
230-
'src_view_name': src_view.get('name'),
231-
'src_headers': src_headers,
188+
'dst_dtable_uuid': dst_dtable_uuid,
189+
'src_table_name': assets.get('src_table_name'),
190+
'src_view_name': assets.get('src_view_name'),
191+
'src_columns': assets.get('src_columns'),
192+
'src_version': assets.get('src_version'),
232193
'dst_table_id': dst_table_id,
233-
'dst_table_name': dst_table_name,
234-
'dst_headers': dst_headers,
235-
'dst_columns': dst_columns,
236-
'dst_rows': dst_rows,
237-
'lang': 'en' # TODO: lang
194+
'dst_table_name': assets.get('dst_table_name'),
195+
'dst_columns': assets.get('dst_columns'),
196+
'operator': 'dtable-events',
197+
'lang': 'en', # TODO: lang
198+
'dataset_id': dataset_id
238199
})
239200
except Exception as e:
240-
logger.error('sync common dataset error: %s', e)
201+
logging.error('sync common dataset error: %s', e)
241202
continue
242203
else:
243204
if result.get('error_msg'):
244-
logger.error(result['error_msg'])
205+
logging.error(result['error_msg'])
245206
if result.get('error_type') == 'generate_synced_columns_error':
246207
set_common_dataset_sync_invalid(dataset_sync_id, db_session)
247208
continue
248209

249-
dataset_update_map[dataset_sync_id] = dtable_src_version
210+
dataset_update_map[dataset_sync_id] = assets.get('src_version')
250211
sync_count += 1
251212

252213
if sync_count == 1000:
253214
try:
254215
update_sync_time_and_version(db_session, dataset_update_map)
255216
except Exception as e:
256-
logger.error(f'update sync time and src_version failed, error: {e}')
217+
logging.error(f'update sync time and src_version failed, error: {e}')
257218
dataset_update_map = {}
258219
sync_count = 0
259220

260221
if dataset_update_map:
261222
try:
262223
update_sync_time_and_version(db_session, dataset_update_map)
263224
except Exception as e:
264-
logger.error(f'update sync time and src_version failed, error: {e}')
225+
logging.error(f'update sync time and src_version failed, error: {e}')
265226

266227

267228
class CommonDatasetSyncerTimer(Thread):
@@ -279,7 +240,7 @@ def timed_job():
279240
try:
280241
check_common_dataset(db_session)
281242
except Exception as e:
282-
logger.exception('check periodcal common dataset syncs error: %s', e)
243+
logging.exception('check periodcal common dataset syncs error: %s', e)
283244
finally:
284245
db_session.close()
285246

0 commit comments

Comments
 (0)