Skip to main content

How to add task list to existing DAG

I'm trying to automate the creation of batch pipelines from backend mySQLL to BiqQuery and I've ran into issue with my script.

It initialises a DAG with it's own class (DagCreator) and then passes it as a instance variable in my MySQLBatchPipelines class.

The tasks and the order of the tasks are contained in methods inside my MySQLBatchPipelines instance.

How would I add those tasks and the tasks list to my existing DAG?

Here's some example code so you understand my query:

class DagCreator:

    def __init__(
        self,
        dag_id,
        start_date,
        time_delay: int
    ):
        self.dag_id = dag_id,
        self.start_date = start_date,
        self.time_delay = time_delay

        print(self.dag_id)

    def DagToReturn(self):
        args = {
            'owner': 'Data Lake',
            'depends_on_past': True,
            'start_date': self.start_date,
            'email_on_failure': False,
            'email_on_retry': False,
            'on_failure_callback': 'failed_dag_slack_alert',
            'retries': 0,
            'max_active_runs': 1,
            'max_file_size': int(50e6)
        }
        dag_to_return = DAG(
            dag_id=f'{self.dag_id}_batch',
            schedule_interval=f'{self.time_delay} * * *',
            concurrency=4,
            default_args=args
        )
        return dag_to_return

dag = DagCreator('table_name', dt.datetime(2022, 1, 15), 7)
dag_returned = dag.DagToReturn()
dag_returned

This is the second class and it's methods:


class MySQLBatchPipeline:
    '''
    Class to generate MySQL batch pipelines that store CSV's 
    in GCS then import them into BigQuery
    '''
    export_format='CSV' 

    def __init__(
        self,
        dag,
    ....

    # lots of methods that return tasks like these
    def uk_source_table_extract_task(self):
        uk_source_table_extract_task = MySqlToGoogleCloudStorageOperator(
            sql=self.queries['new_rows_batch_query'], # <--- needs testing
            bucket=self.gcs_bucket,
            filename=self.extract_gcs_path,
            approx_max_file_size_bytes=MAX_FILE_SIZE, # <-- is max_file_size necc? is in dag arg's is it necc here
            mysql_conn_id=self.mysql_connection_id,
            export_format=self.export_format,
            google_cloud_storage_conn_id=self.gcs_connection_id,
            params={
                'export_format': self.export_format,
                'country': 'uk',
                'time_delay': self.time_delay
            },
            task_id='uk_source_table_extract_task',
            dag=self.dag
      return uk_expected_stg_row_count_extract_task


    # this the the workflow I want to add to the existing DAG I've passed a instance attribute
        def uk_workflow(self):
            uk_source_table_extract_task.set_downstream(
                [
                    uk_table_name_load_task,
                    uk_truncate_table_name_staging_task,
                    uk_expected_stg_row_count_extract_task
                ]
            )
            uk_truncate_table_name_staging_task.set_downstream(
            uk_table_name_load_task
            # etc etc until the final task


source https://stackoverflow.com/questions/70756217/how-to-add-task-list-to-existing-dag

Comments

Popular posts from this blog

Confusion between commands.Bot and discord.Client | Which one should I use?

Whenever you look at YouTube tutorials or code from this website there is a real variation. Some developers use client = discord.Client(intents=intents) while the others use bot = commands.Bot(command_prefix="something", intents=intents) . Now I know slightly about the difference but I get errors from different places from my code when I use either of them and its confusing. Especially since there has a few changes over the years in discord.py it is hard to find the real difference. I tried sticking to discord.Client then I found that there are more features in commands.Bot . Then I found errors when using commands.Bot . An example of this is: When I try to use commands.Bot client = commands.Bot(command_prefix=">",intents=intents) async def load(): for filename in os.listdir("./Cogs"): if filename.endswith(".py"): client.load_extension(f"Cogs.{filename[:-3]}") The above doesnt giveany response from my Cogs ...

How to show number of registered users in Laravel based on usertype?

i'm trying to display data from the database in the admin dashboard i used this: <?php use Illuminate\Support\Facades\DB; $users = DB::table('users')->count(); echo $users; ?> and i have successfully get the correct data from the database but what if i want to display a specific data for example in this user table there is "usertype" that specify if the user is normal user or admin i want to user the same code above but to display a specific usertype i tried this: <?php use Illuminate\Support\Facades\DB; $users = DB::table('users')->count()->WHERE usertype =admin; echo $users; ?> but it didn't work, what am i doing wrong? source https://stackoverflow.com/questions/68199726/how-to-show-number-of-registered-users-in-laravel-based-on-usertype

Why is my reports service not connecting?

I am trying to pull some data from a Postgres database using Node.js and node-postures but I can't figure out why my service isn't connecting. my routes/index.js file: const express = require('express'); const router = express.Router(); const ordersCountController = require('../controllers/ordersCountController'); const ordersController = require('../controllers/ordersController'); const weeklyReportsController = require('../controllers/weeklyReportsController'); router.get('/orders_count', ordersCountController); router.get('/orders', ordersController); router.get('/weekly_reports', weeklyReportsController); module.exports = router; My controllers/weeklyReportsController.js file: const weeklyReportsService = require('../services/weeklyReportsService'); const weeklyReportsController = async (req, res) => { try { const data = await weeklyReportsService; res.json({data}) console...