/
MaxIs
/
savemoney
Обзор
Документация
Войти
/
MaxIs
/
savemoney
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
sm/api/tasks.py
240 строк
10 KB
zalma2006
добавлен flask для работы с selenium в отдельном сервисе, добавлены точки API, добавлена загрузка отчётов из источника, подробнее в HISTORY.md
15 мар 2026, 19:36
15 мар 2026, 19:36
ae3ea09
Код
Авторство
О чём код?
import datetime import os from pathlib import Path import pandas as pd import requests from celery import chain, group from logging import getLogger from savemoney.models import EmitentReport from sm.celery import app from savemoney.data_spravochniki import REPORT_COLS_TYPE from utils.cust_func import get_simple_json2resp logger = getLogger(__name__) @app.task(bind=True) def download_file_url(self, url: str, schema: str): furl = os.getenv('SMD_URL', None) if not furl: err = 'SMD_URL is not set' logger.error(err) raise Exception(err) furl = f'{schema}://{furl}/download_file' headers = { 'x-api-key': os.getenv('SMD_API_KEY'), 'Content-Type': 'application/json', } send_data = { 'url': url, } data = get_simple_json2resp() resp = requests.post(url=furl, headers=headers, json=send_data) resp.raise_for_status() resp_json = resp.json() if len(resp_json['errors']) == 0: data['tmp_path'] = resp_json.get('filename', None) data['message'] = 'success' data['download_time_sec'] = resp_json.get('download_time_sec', None) else: data['errors'] = resp_json['errors'] data['url'] = url return data @app.task(bind=True) def upload_report(self, req, send_data, schema): if len(req['errors']) == 0: filepath = Path(req['tmp_path']) download_time_sec = float(req['download_time_sec']) url_path = req['url'] send_data.update({'filepath': str(filepath), 'name': filepath.name, 'download_time_sec': download_time_sec, 'url_path': url_path, }) for d in ['approve_date', 'basis_date', 'loc_date']: date = send_data[d] if date is not None: date = datetime.datetime.strftime(send_data[d], '%d.%m.%Y') send_data[d] = date django_host = os.getenv('DJANGO_HOST', '127.0.0.1') django_port = os.getenv('DJANGO_PORT', 8000) api_key = os.getenv('API_KEY', None) url = f'{schema}://{django_host}:{django_port}/api/upload_report/path' headers = {'x-api-key': api_key} request = requests.post(url, headers=headers, data=send_data) if request.status_code >= 400: print(request.content[:500]) request.raise_for_status() else: raise Exception(req['errors']) @app.task(bind=True) def find_report_celery(self, inn: str, r_types: str, idx_company: int, schema: str, *args, **kwargs): """Получить данные по report""" furl = os.getenv('SMD_URL', None) if not furl: err = 'SMD_URL is not set' logger.error(err) raise Exception(err) res = requests.post(f"{schema}://{furl}/find_report", json={'schema': schema, 'inn': inn, 'r_types': r_types, 'idx_company': idx_company}, headers={'x-api-key': os.getenv('SMD_API_KEY'), 'Content-Type': 'application/json'}) if res.ok: data = res.json() data['message'] = 'success' else: try: err = res.json() err = err.get('errors', ['no found errors in find_report_celery!']) except Exception: err = [res.content] data = { 'errors': err, 'message': '', } data['status_code'] = res.status_code return data @app.task(bind=True) def set_idx_source_emitent(self, data, inn: str, idx_company: int, schema: str, source_url: str): """Создаём связь ресурс - эмитент""" if len(data['errors']) == 0: # создаём связь Resource <-> Emitent if idx_company is None and data.get('idx_company', None) is not None: idx_company = data.get('idx_company', None) django_host = os.getenv('DJANGO_HOST', '127.0.0.1') django_port = os.getenv('DJANGO_PORT', 8000) api_key = os.getenv('API_KEY', None) url = f'{schema}://{django_host}:{django_port}/api/resource_emitent/add' headers = {'x-api-key': api_key, 'Content-Type': 'application/json; charset=utf-8'} send_data = { 'inn': inn, 'source': source_url, 'source_sv': idx_company, } cr_res_em = requests.post(url, headers=headers, json=send_data) if cr_res_em.ok: cr_res_em = cr_res_em.json() if len(cr_res_em['errors']) == 0: data['message'] += f" {cr_res_em['message']}" else: data['errors'].append(cr_res_em['errors']) else: data['errors'].append(f" not success create Resource <-> Emitent, " f"status_code = {cr_res_em.status_code}," f" content[:100] = {cr_res_em.content[:100]}") return data @app.task(bind=True) def set_report_emitent(self, res, inn: str, schema: str): """Установка отчётов эмитента""" task_id = self.request.id req = get_simple_json2resp() if len(res['errors']) == 0 and 'reports' in res and len(res['reports']) > 0: req['reports'] = dict() for report_type, content in res['reports'].items(): (need_cols, nd, link, loc_date, emre, basis_date, approve_date, rep_year, source_file, type_doc, rep_period) = [None for _ in range(11)] if str(report_type).isdigit(): report_type = int(report_type) if report_type in REPORT_COLS_TYPE: need_cols = REPORT_COLS_TYPE[int(report_type)].copy() else: req['errors'].append(f"{report_type} is not contained in REPORT_COLS_TYPE!") continue if need_cols is not None: nd = content['need_doc'] if nd is not None: type_doc = nd.get('тип_документа', None) rep_year = nd.get('отч_год', None) rep_period = nd.get('отч_месяц', None) approve_date = nd.get('дата_утверждения', None) basis_date = nd.get('дата_основания', None) loc_date = nd.get('дата_размещения', None) link = nd.get('файл', None) else: req['errors'].append(f"{inn}, report_type = {report_type}, not contained need_doc") continue if link is not None: # получаем последний связанный отчёт emre = EmitentReport.objects.filter( report_type__num_type=report_type, emitent__emitent_inn=inn, report__current_report=True).order_by('-report__loc_date', '-report__basis_date').first() else: req['errors'].append(f"{inn}, report_type = {report_type}, not contained file link") if len(req['errors']) > 0: req['errors'] = '||'.join(req['errors']) return req loc_date = pd.to_datetime(loc_date, yearfirst=True) basis_date = pd.to_datetime(basis_date, yearfirst=True) approve_date = pd.to_datetime(approve_date, yearfirst=True) loc_date = None if pd.isna(loc_date) else loc_date.to_pydatetime().date() basis_date = None if pd.isna(basis_date) else basis_date.to_pydatetime().date() approve_date = None if pd.isna(approve_date) else approve_date.to_pydatetime().date() rep_year = int(rep_year) if rep_year is not None and str(rep_year).isdigit() else None send_data = { 'inn': inn, 'report_type': report_type, 'current_report': True, 'type_doc': type_doc, 'rep_year': rep_year, 'rep_period': rep_period, 'approve_date': approve_date, 'basis_date': basis_date, 'loc_date': loc_date, } if emre is not None: # проверим размещение отчётов если они произошли позже тогда нужно старый элемент отметить старым, а # новый записать как текущий base_report_loc_date = emre.report.loc_date # устанавливаем документ как старый if not pd.isna(loc_date) and base_report_loc_date < loc_date: emre.report.current_report = False emre.report.save() thread = chain( download_file_url.s(url=link, schema=schema), upload_report.s(schema=schema, send_data=send_data), ) thread.apply_async() else: thread = chain( download_file_url.s(url=link, schema=schema), upload_report.s(schema=schema, send_data=send_data), ) thread.apply_async() else: if len(res['errors']) > 0: err = res['errors'] else: err = [] req['errors'].extend(err) if len(req['errors']) == 0: req['message'] = f'success {inn} reports add' return req @app.task(bind=True) def set_report_celery(self, inn: str, r_types: str, idx_company: int, schema: str, source_url: str, *args, **kwargs): """Найти отчёты по эмитентам""" task_id = self.request.id thread = chain( find_report_celery.s(inn, r_types, idx_company, schema, *args, **kwargs), group( set_idx_source_emitent.s(inn=inn, idx_company=idx_company, schema=schema, source_url=source_url), set_report_emitent.s(inn=inn, schema=schema), ), ) res = thread.apply_async() return res