Skip to content

fix: speed up launch-plan loading in large projects - #940

Open
1fanwang wants to merge 1 commit into
flyteorg:masterfrom
1fanwang:1fanwang/fix-launch-plan-query-scope
Open

1fanwang wants to merge 1 commit into
flyteorg:masterfrom
1fanwang:1fanwang/fix-launch-plan-query-scope

Conversation

@1fanwang

Copy link
Copy Markdown
Contributor

TL;DR

The launch dialog and schedules panel can load slowly in projects with many workflows while their launch-plan lookups scan unrelated workflows. Both panels now use the selected workflow's full scope for an indexed lookup instead of waiting for that broader scan.

Type

  • Bug Fix
  • Feature
  • Plugin

Are all requirements met?

  • Code completed
  • Smoke tested
  • Unit tests added
  • Code documentation added
  • Any pending items have an associated Issue

No public documentation needs changing, and no pending items remain.

Complete description

Add the workflow-side scope to both queries. Keep the existing launch-plan scope, version and active-state filters, ordering, limits, and preferred launch-plan lookup unchanged.

Tracking Issue

Related to flyteorg/flyte#6428. This covers two call sites outside the earlier fix at #891.

Follow-up issue

None.

Testing Done

The real React components ran in JSDOM through their HTTP client, public FlyteAdmin, and PostgreSQL 14.19. The catalog was registered through Admin after its normal migrations, not inserted directly into the database.

# Scenario Command Result
1 Launch dialog and schedules Live component probe below Both queries lacked workflow scope before; both include it after.
2 Query access paths psql scripts below Launch lookup uses 3 workflow buffers instead of 31; schedules avoid discarding 3,999 workflows.
3 Launch-dialog result preservation curl and diff below The launch-dialog lookup returns HTTP 200 before and after with identical complete responses. An invalid filter returns HTTP 400.

This is a local PostgreSQL check, not a MySQL benchmark or a full-browser test. The original launch-dialog query already used an index, but without its leading project/domain predicates. Only the original schedule query used a sequential workflow scan in this fixture.

Setup

The backend was flyteorg/flyte@b2be54c, with PostgreSQL on 127.0.0.1:15447, database flytemining, user flyteprobe, and Admin on 127.0.0.1:18097. Admin used memory object storage. A test-only build overlay restricted listener addresses to loopback; it did not change request handling or SQL.

With that migrated local Admin running, register the project:

curl --noproxy '*' --fail-with-body --silent --show-error \
  -H 'Content-Type: application/json' \
  -d '{"project":{"id":"query-probe","name":"query-probe"}}' \
  http://127.0.0.1:18097/api/v1/projects

Save this registration script as seed_catalog.py:

from __future__ import annotations

import json
import sys
import time
from http import HTTPStatus
from typing import Any
from urllib.error import HTTPError
from urllib.request import Request, urlopen


def call(*, method: str, path: str, body: dict[str, Any]) -> None:
    request = Request(
        url=f"http://127.0.0.1:18097/api/v1/{path}",
        data=json.dumps(body).encode(),
        headers={"Content-Type": "application/json"},
        method=method,
    )
    try:
        with urlopen(request, timeout=30) as response:
            assert response.status == HTTPStatus.OK, response.status
            response.read()
    except HTTPError as error:
        print(error.read().decode(), file=sys.stderr)
        raise


def main() -> None:
    count = int(sys.argv[1])
    assert count > 0
    scope = {"project": "query-probe", "domain": "development", "version": "v1"}
    task_id = {**scope, "resource_type": "TASK", "name": "noop"}
    call(
        method="POST",
        path="tasks",
        body={
            "id": task_id,
            "spec": {
                "template": {
                    "id": task_id,
                    "type": "python-task",
                    "interface": {"inputs": {}, "outputs": {}},
                    "metadata": {
                        "runtime": {
                            "type": "FLYTE_SDK",
                            "version": "1.0.0",
                            "flavor": "python",
                        }
                    },
                    "container": {
                        "image": "python:3.12",
                        "command": ["python"],
                        "args": ["-c", "pass"],
                    },
                }
            },
        },
    )
    started = time.monotonic()
    for index in range(count):
        name = f"workflow-{index:05d}"
        workflow_id = {**scope, "resource_type": "WORKFLOW", "name": name}
        call(
            method="POST",
            path="workflows",
            body={
                "id": workflow_id,
                "spec": {
                    "template": {
                        "id": workflow_id,
                        "interface": {"inputs": {}, "outputs": {}},
                        "nodes": [
                            {"id": "n0", "task_node": {"reference_id": task_id}}
                        ],
                    }
                },
            },
        )
        launch_plan_id = {**scope, "resource_type": "LAUNCH_PLAN", "name": name}
        call(
            method="POST",
            path="launch_plans",
            body={
                "id": launch_plan_id,
                "spec": {
                    "workflow_id": workflow_id,
                    "default_inputs": {},
                    "fixed_inputs": {},
                    "entity_metadata": {},
                },
            },
        )
        if (index + 1) % 1000 == 0:
            print(f"Registered {index + 1} workflows and launch plans", flush=True)
    call(
        method="PUT",
        path=f"launch_plans/query-probe/development/{name}/v1",
        body={"id": launch_plan_id, "state": "ACTIVE"},
    )
    print(f"Catalog registration finished in {time.monotonic() - started:.2f}s")
    print(f"Selected workflow: query-probe/development/{name}/v1")


if __name__ == "__main__":
    main()
python3 seed_catalog.py 4000
psql -h 127.0.0.1 -p 15447 -U flyteprobe -d flytemining \
  -v ON_ERROR_STOP=1 -c 'ANALYZE workflows' -c 'ANALYZE launch_plans'
Registered 1000 workflows and launch plans
Registered 2000 workflows and launch plans
Registered 3000 workflows and launch plans
Registered 4000 workflows and launch plans
Catalog registration finished in 10.39s
Selected workflow: query-probe/development/workflow-03999/v1
ANALYZE
ANALYZE

Live component requests

Save this temporary probe as packages/oss-console/src/components/Launch/LaunchForm/test/LaunchPlanQueries.live.test.tsx. It restores the real HTTP transport that the ordinary test setup stubs; the observation spy still sends every request to Admin.

import React from 'react';
import { render, screen, waitFor } from '@testing-library/react';
import { ThemeProvider } from '@mui/material/styles';
import { muiTheme } from '@clients/theme/Theme/muiTheme';
import { QueryClient, QueryClientProvider } from 'react-query';
import { MemoryRouter } from 'react-router';
import { ResourceIdentifier, ResourceType } from '../../../../models/Common/types';
import { EntitySchedules } from '../../../Entities/EntitySchedules';
import { WorkflowNodeExecutionsProvider } from '../../../Executions/contextProvider/NodeExecutionDetails/WorkflowNodeExecutionsProvider';
import { LaunchForm } from '../LaunchForm';

const { fetch: realFetch } = jest.requireActual<typeof import('whatwg-fetch')>('whatwg-fetch');

jest.setTimeout(20000);

const id: ResourceIdentifier = {
  resourceType: ResourceType.WORKFLOW,
  project: 'query-probe',
  domain: 'development',
  name: 'workflow-03999',
};

describe('Launch-plan queries against public Admin and PostgreSQL', () => {
  let client: QueryClient;
  let requests: jest.SpyInstance;

  beforeEach(() => {
    client = new QueryClient({
      defaultOptions: { queries: { retry: false, cacheTime: 0 } },
    });
    requests = jest.spyOn(globalThis, 'fetch').mockImplementation(realFetch);
  });

  afterEach(() => {
    client.clear();
    requests.mockRestore();
  });

  const renderLive = (element: React.ReactElement) =>
    render(
      <ThemeProvider theme={muiTheme}>
        <QueryClientProvider client={client}>
          <MemoryRouter>
            <WorkflowNodeExecutionsProvider initialNodeExecutions={[]}>
              {element}
            </WorkflowNodeExecutionsProvider>
          </MemoryRouter>
        </QueryClientProvider>
      </ThemeProvider>,
    );

  const assertWorkflowScope = (surface: string) => {
    const request = requests.mock.calls.find(([url]) =>
      String(url).includes('/api/v1/launch_plans/'),
    );
    if (!request) {
      throw new Error(`${surface} did not request launch plans`);
    }
    const url = new URL(String(request[0]));
    const filters = url.searchParams.get('filters');
    console.log(`${surface} returned ${id.name}; filters=${filters}`);
    expect(filters).toContain(`eq(workflow.name,${id.name})`);
    expect(filters).toContain(`eq(workflow.project,${id.project})`);
    expect(filters).toContain(`eq(workflow.domain,${id.domain})`);
  };

  it('scopes the launch form query to the selected workflow', async () => {
    const { container } = renderLive(<LaunchForm workflowId={id} onClose={() => {}} />);
    await waitFor(
      () => expect(container.querySelector('#launch-lp-selector')).toHaveValue(id.name),
      { timeout: 15000 },
    );
    assertWorkflowScope('Launch form');
  });

  it('scopes the schedule query to the selected workflow', async () => {
    renderLive(<EntitySchedules id={id} />);
    await screen.findByText('Schedules', {}, { timeout: 15000 });
    assertWorkflowScope('Schedules');
  });
});

Run from the console checkout with its locked dependencies installed:

ADMIN_API=http://127.0.0.1:18097 node node_modules/jest/bin/jest.js \
  --runInBand --runTestsByPath \
  packages/oss-console/src/components/Launch/LaunchForm/test/LaunchPlanQueries.live.test.tsx

Before, at 04d6a33:

Launch form returned workflow-03999; filters=eq(workflow.name,workflow-03999)+eq(workflow.version,v1)
Schedules returned workflow-03999; filters=eq(state,1)+eq(workflow.name,workflow-03999)
Expected substring: "eq(workflow.project,query-probe)"
Received string:    "eq(workflow.name,workflow-03999)+eq(workflow.version,v1)"
Expected substring: "eq(workflow.project,query-probe)"
Received string:    "eq(state,1)+eq(workflow.name,workflow-03999)"

After:

Launch form returned workflow-03999; filters=eq(workflow.name,workflow-03999)+eq(workflow.version,v1)+eq(workflow.project,query-probe)+eq(workflow.domain,development)
Schedules returned workflow-03999; filters=eq(state,1)+eq(workflow.name,workflow-03999)+eq(workflow.project,query-probe)+eq(workflow.domain,development)

Query plans

These are the queries captured from the two component requests. Save the original statements as public/query-plans-before.sql:

-- Exact unpatched UI queries captured from public FlyteAdmin PostgreSQL logs.
PREPARE probe_before_0 AS SELECT "launch_plans"."id","launch_plans"."created_at","launch_plans"."updated_at","launch_plans"."deleted_at","launch_plans"."project","launch_plans"."domain","launch_plans"."name","launch_plans"."version","launch_plans"."spec","launch_plans"."workflow_id","launch_plans"."closure","launch_plans"."state","launch_plans"."digest","launch_plans"."schedule_type","launch_plans"."launch_condition_type" FROM "launch_plans" inner join workflows on launch_plans.workflow_id = workflows.id WHERE launch_plans.project = $1 AND launch_plans.domain = $2 AND workflows.name = $3 AND workflows.version = $4 LIMIT 50;
EXPLAIN (ANALYZE, BUFFERS) EXECUTE probe_before_0('query-probe', 'development', 'workflow-03999', 'v1');
PREPARE probe_before_1 AS SELECT "launch_plans"."id","launch_plans"."created_at","launch_plans"."updated_at","launch_plans"."deleted_at","launch_plans"."project","launch_plans"."domain","launch_plans"."name","launch_plans"."version","launch_plans"."spec","launch_plans"."workflow_id","launch_plans"."closure","launch_plans"."state","launch_plans"."digest","launch_plans"."schedule_type","launch_plans"."launch_condition_type" FROM "launch_plans" inner join workflows on launch_plans.workflow_id = workflows.id WHERE launch_plans.project = $1 AND launch_plans.domain = $2 AND launch_plans.state = $3 AND workflows.name = $4 ORDER BY created_at desc LIMIT 5;
EXPLAIN (ANALYZE, BUFFERS) EXECUTE probe_before_1('query-probe', 'development', '1', 'workflow-03999');

Save the scoped statements as public/query-plans-after.sql:

-- Exact patched UI queries captured from public FlyteAdmin PostgreSQL logs.
PREPARE probe_after_0 AS SELECT "launch_plans"."id","launch_plans"."created_at","launch_plans"."updated_at","launch_plans"."deleted_at","launch_plans"."project","launch_plans"."domain","launch_plans"."name","launch_plans"."version","launch_plans"."spec","launch_plans"."workflow_id","launch_plans"."closure","launch_plans"."state","launch_plans"."digest","launch_plans"."schedule_type","launch_plans"."launch_condition_type" FROM "launch_plans" inner join workflows on launch_plans.workflow_id = workflows.id WHERE launch_plans.project = $1 AND launch_plans.domain = $2 AND workflows.name = $3 AND workflows.version = $4 AND workflows.project = $5 AND workflows.domain = $6 LIMIT 50;
EXPLAIN (ANALYZE, BUFFERS) EXECUTE probe_after_0('query-probe', 'development', 'workflow-03999', 'v1', 'query-probe', 'development');
PREPARE probe_after_1 AS SELECT "launch_plans"."id","launch_plans"."created_at","launch_plans"."updated_at","launch_plans"."deleted_at","launch_plans"."project","launch_plans"."domain","launch_plans"."name","launch_plans"."version","launch_plans"."spec","launch_plans"."workflow_id","launch_plans"."closure","launch_plans"."state","launch_plans"."digest","launch_plans"."schedule_type","launch_plans"."launch_condition_type" FROM "launch_plans" inner join workflows on launch_plans.workflow_id = workflows.id WHERE launch_plans.project = $1 AND launch_plans.domain = $2 AND launch_plans.state = $3 AND workflows.name = $4 AND workflows.project = $5 AND workflows.domain = $6 ORDER BY created_at desc LIMIT 5;
EXPLAIN (ANALYZE, BUFFERS) EXECUTE probe_after_1('query-probe', 'development', '1', 'workflow-03999', 'query-probe', 'development');
psql -h 127.0.0.1 -p 15447 -U flyteprobe -d flytemining \
  -v ON_ERROR_STOP=1 -P pager=off -f public/query-plans-before.sql
psql -h 127.0.0.1 -p 15447 -U flyteprobe -d flytemining \
  -v ON_ERROR_STOP=1 -P pager=off -f public/query-plans-after.sql

Before:

PREPARE
                                                                        QUERY PLAN                                                                        
----------------------------------------------------------------------------------------------------------------------------------------------------------
 Limit  (cost=0.56..166.60 rows=1 width=268) (actual time=0.157..0.157 rows=1 loops=1)
   Buffers: shared hit=34
   ->  Nested Loop  (cost=0.56..166.60 rows=1 width=268) (actual time=0.156..0.157 rows=1 loops=1)
         Buffers: shared hit=34
         ->  Index Scan using workflow_project_domain_name_idx on workflows  (cost=0.28..158.29 rows=1 width=8) (actual time=0.153..0.153 rows=1 loops=1)
               Index Cond: (name = 'workflow-03999'::text)
               Filter: (version = 'v1'::text)
               Buffers: shared hit=31
         ->  Index Scan using idx_launch_plans_workflow_id on launch_plans  (cost=0.28..8.30 rows=1 width=268) (actual time=0.002..0.002 rows=1 loops=1)
               Index Cond: (workflow_id = workflows.id)
               Filter: ((project = 'query-probe'::text) AND (domain = 'development'::text))
               Buffers: shared hit=3
 Planning:
   Buffers: shared hit=437
 Planning Time: 10.584 ms
 Execution Time: 0.177 ms
(16 rows)

PREPARE
                                                                          QUERY PLAN                                                                           
---------------------------------------------------------------------------------------------------------------------------------------------------------------
 Limit  (cost=158.32..158.33 rows=1 width=268) (actual time=0.384..0.385 rows=1 loops=1)
   Buffers: shared hit=106
   ->  Sort  (cost=158.32..158.33 rows=1 width=268) (actual time=0.384..0.385 rows=1 loops=1)
         Sort Key: launch_plans.created_at DESC
         Sort Method: quicksort  Memory: 25kB
         Buffers: shared hit=106
         ->  Nested Loop  (cost=0.28..158.31 rows=1 width=268) (actual time=0.374..0.375 rows=1 loops=1)
               Buffers: shared hit=103
               ->  Seq Scan on workflows  (cost=0.00..150.00 rows=1 width=8) (actual time=0.372..0.372 rows=1 loops=1)
                     Filter: (name = 'workflow-03999'::text)
                     Rows Removed by Filter: 3999
                     Buffers: shared hit=100
               ->  Index Scan using idx_launch_plans_workflow_id on launch_plans  (cost=0.28..8.30 rows=1 width=268) (actual time=0.001..0.002 rows=1 loops=1)
                     Index Cond: (workflow_id = workflows.id)
                     Filter: ((project = 'query-probe'::text) AND (domain = 'development'::text) AND (state = 1))
                     Buffers: shared hit=3
 Planning:
   Buffers: shared hit=23
 Planning Time: 0.095 ms
 Execution Time: 0.391 ms
(20 rows)

After:

PREPARE
                                                                       QUERY PLAN                                                                        
---------------------------------------------------------------------------------------------------------------------------------------------------------
 Limit  (cost=0.56..16.62 rows=1 width=268) (actual time=0.013..0.014 rows=1 loops=1)
   Buffers: shared hit=6
   ->  Nested Loop  (cost=0.56..16.62 rows=1 width=268) (actual time=0.013..0.013 rows=1 loops=1)
         Buffers: shared hit=6
         ->  Index Scan using workflow_project_domain_name_idx on workflows  (cost=0.28..8.30 rows=1 width=8) (actual time=0.009..0.010 rows=1 loops=1)
               Index Cond: ((project = 'query-probe'::text) AND (domain = 'development'::text) AND (name = 'workflow-03999'::text))
               Filter: (version = 'v1'::text)
               Buffers: shared hit=3
         ->  Index Scan using idx_launch_plans_workflow_id on launch_plans  (cost=0.28..8.30 rows=1 width=268) (actual time=0.002..0.003 rows=1 loops=1)
               Index Cond: (workflow_id = workflows.id)
               Filter: ((project = 'query-probe'::text) AND (domain = 'development'::text))
               Buffers: shared hit=3
 Planning:
   Buffers: shared hit=440
 Planning Time: 7.965 ms
 Execution Time: 0.037 ms
(16 rows)

PREPARE
                                                                          QUERY PLAN                                                                           
---------------------------------------------------------------------------------------------------------------------------------------------------------------
 Limit  (cost=16.63..16.63 rows=1 width=268) (actual time=0.015..0.016 rows=1 loops=1)
   Buffers: shared hit=9
   ->  Sort  (cost=16.63..16.63 rows=1 width=268) (actual time=0.015..0.015 rows=1 loops=1)
         Sort Key: launch_plans.created_at DESC
         Sort Method: quicksort  Memory: 25kB
         Buffers: shared hit=9
         ->  Nested Loop  (cost=0.56..16.62 rows=1 width=268) (actual time=0.004..0.005 rows=1 loops=1)
               Buffers: shared hit=6
               ->  Index Scan using workflow_project_domain_name_idx on workflows  (cost=0.28..8.30 rows=1 width=8) (actual time=0.002..0.002 rows=1 loops=1)
                     Index Cond: ((project = 'query-probe'::text) AND (domain = 'development'::text) AND (name = 'workflow-03999'::text))
                     Buffers: shared hit=3
               ->  Index Scan using idx_launch_plans_workflow_id on launch_plans  (cost=0.28..8.30 rows=1 width=268) (actual time=0.002..0.002 rows=1 loops=1)
                     Index Cond: (workflow_id = workflows.id)
                     Filter: ((project = 'query-probe'::text) AND (domain = 'development'::text) AND (state = 1))
                     Buffers: shared hit=3
 Planning:
   Buffers: shared hit=23
 Planning Time: 0.102 ms
 Execution Time: 0.024 ms
(19 rows)

Launch-dialog response and validation controls

curl --noproxy '*' --get --fail-with-body --silent --show-error \
  'http://127.0.0.1:18097/api/v1/launch_plans/query-probe/development' \
  --data-urlencode 'filters=eq(workflow.name,workflow-03999)+eq(workflow.version,v1)' \
  --data-urlencode 'limit=50' -o publish-api-before.json -w 'Before HTTP %{http_code}\n'
curl --noproxy '*' --get --fail-with-body --silent --show-error \
  'http://127.0.0.1:18097/api/v1/launch_plans/query-probe/development' \
  --data-urlencode 'filters=eq(workflow.name,workflow-03999)+eq(workflow.version,v1)+eq(workflow.project,query-probe)+eq(workflow.domain,development)' \
  --data-urlencode 'limit=50' -o publish-api-after.json -w 'After HTTP %{http_code}\n'
diff -u publish-api-before.json publish-api-after.json
jq '{count:(.launchPlans|length),names:[.launchPlans[].id.name]}' publish-api-after.json
curl --noproxy '*' --get --silent --show-error \
  'http://127.0.0.1:18097/api/v1/launch_plans/query-probe/development' \
  --data-urlencode 'filters=eq(workflow.bogusfield,x)' --data-urlencode 'limit=1' \
  -o publish-api-invalid.json -w 'Invalid filter HTTP %{http_code}\n'
cat publish-api-invalid.json
Before HTTP 200
After HTTP 200
{
  "count": 1,
  "names": [
    "workflow-03999"
  ]
}
Invalid filter HTTP 400
{"code":3, "message":"'w.bogusfield' is invalid filter", "details":[]}

The temporary live probe is not part of CI.

  • Local code review completed (no separate review was run).

Signed-off-by: 1fanwang <1fannnw@gmail.com>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant