Но я получаю ошибку при попытке использовать ее в своем даге.
Затем я запускаю свой даг и получаю ошибку для подключения_to_api задача:
{python.py:202} INFO - Done. Returned value was: []
{xcom.py:664} ERROR - Object of type yclientsapi is not JSON serializable. If you are using pickle instead of JSON for XCom, then you need to enable pickle support for XCom in your airflow config or make sure to decorate your object with attr.
TypeError: Object of type yclientsapi is not JSON serializable
Я проверил сериализацию в документации Airflow, но не думаю, что хочу переписывать код библиотеки.
Как мне нужно исправить свой код, чтобы избежать этой ошибки?
П. S. Я попробовал включить поддержку Pickle для XCom, добавив в docker-compose.yml AIRFLOW__CORE__ENABLE_XCOM_PICKLING=true.
Это не решило мою проблему.

Airflow 2.8.2, Python 3.11
Код моей библиотеки:
import pandas as pd
import numpy as np
from datetime import date
import httpx
import ujson
# create class
class yclientsapi:
def __init__(self, bearer_key: str, user_key: str):
self.bearer_key = bearer_key
self.user_key = user_key
self.headers = {
"Accept": "application/vnd.yclients.v2+json",
"Content-Type": "application/json",
'Authorization': f"Bearer {bearer_key}, User {user_key}"
}
# get table with salons (stores)
def get_salons_info(self):
session = httpx.Client()
url = "https://api.yclients.com/api/v1/groups"
response = session.get(url, headers=self.headers)
first_request= ujson.loads(response.text)
df = pd.DataFrame({
'salon_id': [i['id'] for i in first_request['data'][0]['companies']]
, 'salon_name': [i['title'] for i in first_request['data'][0]['companies']]
})
salon_ids_list = [i['id'] for i in first_request['data'][0]['companies']]
return df, salon_ids_list
Код My Dag:
import datetime
import yclients as yc
from airflow.decorators import dag, task
from airflow.models import Variable
import httpx
default_args = {
'owner': 'user',
'depends_on_past': False,
'retries': 2,
'retry_delay': datetime.timedelta(minutes=5),
'start_date': datetime.datetime(2024, 6, 20)
}
schedule_interval = '*/20 * * * *'
host = Variable.get('host')
database_name = Variable.get('database_name')
user_name = Variable.get('user_name')
password_for_db = Variable.get('password_for_db')
server_host_name = Variable.get('server_host_name')
bearer_key = Variable.get('bearer_key')
user_key = Variable.get('user_key')
sales_plans_url = Variable.get('sales_plans_url')
specialization_prices_url = Variable.get('specialization_prices_url')
bot_token = Variable.get('bot_token')
chat_id = Variable.get('chat_id')
@dag(default_args=default_args, schedule_interval=schedule_interval, catchup=False, concurrency=4)
def dag_update_database_test():
@task
def connect_to_api(bearer_key: str, user_key: str):
api = yc.yclientsapi(bearer_key, user_key)
result = []
result.append(api)
return result
@task
def get_salons(api):
api = api[0]
result_from_salons_api = api.get_salons_info()
return result_from_salons_api
@task
def send_msg(bot_token: str, chat_id: str, message: str):
message = message[0]
url = f'https://api.telegram.org/bot{bot_token} ... t={message}'
client = httpx.Client()
client.post(url)
api_connection_task = connect_to_api(bearer_key=bearer_key, user_key=user_key)
# get_salon_task[0] = salons_df, get_salon_task[1] = salon_ids_list
get_salon_task = get_salons(api_connection_task)
send_test_message = send_msg(bot_token, chat_id, get_salon_task)
api_connection_task.set_downstream(get_salon_task)
get_salon_task.set_downstream(send_test_message)
dag_update_database_test = dag_update_database_test()
Подробнее здесь: https://stackoverflow.com/questions/786 ... in-airflow