composerairflow_impl

package
v0.59.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 17, 2026 License: Apache-2.0 Imports: 23 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var AirflowDagProcessorManagerLogFilterTask = componentFilterTask(composerairflow.AirflowDagProcessorManagerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "dag-processor-manager")
View Source
var AirflowDagProcessorManagerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask(
	composerairflow.AirflowDagProcessorManagerLogGrouperTaskID,
	composerairflow.AirflowDagProcessorManagerLogFilterTaskID.Ref(),
	func(ctx context.Context, l *log.Log) string {
		fs, err := composerairflow.ExtractComposer(l.NodeReader)
		if err != nil {
			return ""
		}
		if fs.SchedulerID != "" {
			return fs.SchedulerID
		}
		return fs.DagProcessorManagerID
	},
)

AirflowDagProcessorManagerLogGrouperTask groups Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogIngesterTask = inspectiontaskbase.NewGroupedLogIngesterTask(
	composerairflow.AirflowDagProcessorManagerLogIngesterTaskID,
	&dagProcessorManagerLogIngester{},
)

AirflowDagProcessorManagerLogIngesterTask is the task that ingests Airflow DAG processor manager logs.

View Source
var AirflowDagProcessorManagerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	composerairflow.AirflowDagProcessorManagerLogToTimelineMapperTaskID,
	&dagProcessorManagerTimelineMapper{
		targetLogType: composerairflow.LogTypeManagedAirflowEnvironment,
		dagFilePath:   "/home/airflow/gcs/dags",
	},
)

AirflowDagProcessorManagerLogToTimelineMapperTask is the task that maps Airflow DAG processor manager logs to timeline events.

View Source
var AirflowOtherLogFilterTask = inspectiontaskbase.NewLogFilterTask(
	composerairflow.AirflowOtherLogFilterTaskID,
	composerairflow.ComposerLogsQueryTaskID.Ref(),
	func(ctx context.Context, l *log.Log) bool {
		component, err := composerairflow.ExtractComposerComponent(l.NodeReader)
		if err != nil {
			return false
		}

		return component != "airflow-worker" && component != "airflow-scheduler" && component != "dag-processor-manager"
	},
)

AirflowOtherLogGrouperTask groups other Airflow logs.

View Source
var AirflowOtherLogIngesterTask = inspectiontaskbase.NewLogIngesterTask(
	composerairflow.AirflowOtherLogIngesterTaskID,
	&otherLogIngester{},
)

AirflowOtherLogIngesterTask is the task that ingests other Airflow logs.

View Source
var AirflowOtherLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	composerairflow.AirflowOtherLogToTimelineMapperTaskID,
	&otherLogToTimelineMapper{},
)

AirflowOtherLogToTimelineMapperTask is the task that maps other Airflow logs to timeline events.

View Source
var AirflowSchedulerLogFilterTask = componentFilterTask(composerairflow.AirflowSchedulerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "airflow-scheduler")

AirflowSchedulerLogGrouperTask groups Airflow scheduler logs.

View Source
var AirflowSchedulerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask(
	composerairflow.AirflowSchedulerLogIngesterTaskID,
	&schedulerLogIngester{},
)

AirflowSchedulerLogIngesterTask is the task that ingests Airflow scheduler logs.

View Source
var AirflowSchedulerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	composerairflow.AirflowSchedulerLogToTimelineMapperTaskID,
	&schedulerLogToTimelineMapper{},
)

AirflowSchedulerLogToTimelineMapperTask is the task that maps Airflow scheduler logs to timeline events.

View Source
var AirflowWorkerLogFilterTask = componentFilterTask(composerairflow.AirflowWorkerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "airflow-worker")

AirflowWorkerLogGrouperTask groups Airflow worker logs.

View Source
var AirflowWorkerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask(
	composerairflow.AirflowWorkerLogIngesterTaskID,
	&workerLogIngester{},
)

AirflowWorkerLogIngesterTask is the task that ingests Airflow worker logs.

View Source
var AirflowWorkerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask(
	composerairflow.AirflowWorkerLogToTimelineMapperTaskID,
	&workerLogToTimelineMapper{},
)

AirflowWorkerLogToTimelineMapperTask is the task that maps Airflow worker logs to timeline events.

View Source
var AutocompleteComposerComponentsTask = inspectiontaskbase.NewGlobalCachedTask(composerairflow.AutocompleteComposerComponentsTaskID, []coretask.Dependency{
	composercluster.ClusterIdentityTaskID.Ref(),
	gcpcommon.InputStartTimeTaskID.Ref(),
	gcpcommon.InputEndTimeTaskID.Ref(),
	composercluster.InputComposerEnvironmentNameTaskID.Ref(),
	gcpcommon.APIClientFactoryTaskID.Ref(),
	gcpcommon.APIClientCallOptionsInjectorTaskID.Ref(),
}, func(ctx context.Context, prevValue inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]) (inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]], error) {
	clusterIdentity := coretask.GetTaskResult(ctx, composercluster.ClusterIdentityTaskID.Ref())
	projectID := clusterIdentity.ProjectID
	location := clusterIdentity.Location

	startTime := coretask.GetTaskResult(ctx, gcpcommon.InputStartTimeTaskID.Ref())
	endTime := coretask.GetTaskResult(ctx, gcpcommon.InputEndTimeTaskID.Ref())
	environmentName := coretask.GetTaskResult(ctx, composercluster.InputComposerEnvironmentNameTaskID.Ref())
	cf := coretask.GetTaskResult(ctx, gcpcommon.APIClientFactoryTaskID.Ref())
	optionInjector := coretask.GetTaskResult(ctx, gcpcommon.APIClientCallOptionsInjectorTaskID.Ref())

	currentDigest := fmt.Sprintf("%s-%s-%s-%s-%d-%d", projectID, location, environmentName, "logging.googleapis.com/log_entry_count", startTime.Unix(), endTime.Unix())
	if currentDigest == prevValue.DependencyDigest {
		return prevValue, nil
	}

	if projectID == "" || environmentName == "" || location == "" {
		return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
			Value: &inspectioncore.AutocompleteResult[string]{
				Values: []string{},
				Hint:   "Components are suggested after the project ID, location, and environment name are provided.",
			},
			DependencyDigest: currentDigest,
		}, nil
	}

	client, err := cf.MonitoringMetricClient(ctx, googlecloud.Project(projectID))
	if err != nil {
		return prevValue, fmt.Errorf("failed to create monitoring metric client: %w", err)
	}

	ctx = optionInjector.InjectToCallContext(ctx, googlecloud.Project(projectID))

	filter := fmt.Sprintf(`resource.type = "cloud_composer_environment" AND metric.type = "logging.googleapis.com/log_entry_count" AND resource.labels.environment_name = "%s" AND resource.labels.location = "%s"`, environmentName, location)

	errorString := ""
	hintString := ""
	metricsLabels, err := googlecloud.QueryResourceLabelsFromMetrics(ctx, client, projectID, filter, startTime, endTime, []string{"metric.label.log"})
	if err != nil {
		errorString = err.Error()
	}

	componentsMap := make(map[string]struct{})
	for _, labels := range metricsLabels {
		if logName, ok := labels["log"]; ok && logName != "" {
			componentsMap[logName] = struct{}{}
		}
	}

	components := make([]string, 0, len(componentsMap))
	for comp := range componentsMap {
		components = append(components, comp)
	}
	sort.Strings(components)

	if hintString == "" && errorString == "" && len(components) == 0 {
		hintString = "No components found for the specified environment and time range."
	}

	return inspectiontaskbase.CacheableTaskResult[*inspectioncore.AutocompleteResult[string]]{
		DependencyDigest: currentDigest,
		Value: &inspectioncore.AutocompleteResult[string]{
			Values: components,
			Error:  errorString,
			Hint:   hintString,
		},
	}, nil
})
View Source
var ComposerLogsQueryTask = gcpcommon.NewStructuredListLogEntriesTask(&composerListLogEntriesTaskSetting{
	taskId:    composerairflow.ComposerLogsQueryTaskID,
	queryName: "Composer Environment Logs",
})

ComposerLogsQueryTask defines a task that gathers logs from Cloud Logging for multiple Composer components.

View Source
var ComposerLogsTailTask = coretask.NewTailTask(
	composerairflow.ComposerLogsTailTaskID,
	[]coretask.Dependency{
		composerairflow.AirflowWorkerLogToTimelineMapperTaskID.Ref(),
		composerairflow.AirflowSchedulerLogToTimelineMapperTaskID.Ref(),
		composerairflow.AirflowDagProcessorManagerLogToTimelineMapperTaskID.Ref(),
		composerairflow.AirflowOtherLogToTimelineMapperTaskID.Ref(),
	},
	inspectioncore.FeatureTaskLabel(
		"Managed Service for Apache Airflow Logs",
		"Gather Managed Service for Apache Airflow logs, including airflow-worker, airflow-scheduler, and airflow-dag-processor-manager, to visualize general environment operations on timelines.",
		2500,
		true,
	),
)
View Source
var InputComposerComponentsTask = formtask.NewSetFormTaskBuilder(composerairflow.InputComposerComponentsTaskID, gcpcommon.FormBasePriority+3000, "Composer Components").
	WithDependencies([]coretask.Dependency{composerairflow.AutocompleteComposerComponentsTaskID.Ref()}).
	WithDefaultValueConstant([]string{"@any"}, true).
	WithAllowAddAll(false).
	WithAllowRemoveAll(false).
	WithAllowCustomValue(false).
	WithDescription(`Select which Composer V3 components to fetch logs from.`).
	WithOptionsFunc(func(ctx context.Context, previousValues []string) ([]inspectionmetadata.SetParameterFormFieldOptionItem, error) {
		autocompleteResult := coretask.GetTaskResult(ctx, composerairflow.AutocompleteComposerComponentsTaskID.Ref())

		var options []inspectionmetadata.SetParameterFormFieldOptionItem
		options = append(options, inspectionmetadata.SetParameterFormFieldOptionItem{
			ID: "@any",
		})
		if autocompleteResult != nil {
			for _, comp := range autocompleteResult.Values {
				options = append(options, inspectionmetadata.SetParameterFormFieldOptionItem{
					ID: comp,
				})
			}
		}
		return options, nil
	}).
	WithHintFunc(func(ctx context.Context, value []string, convertedValue any) (string, inspectionmetadata.ParameterHintType, error) {
		autocompleteResult := coretask.GetTaskResult(ctx, composerairflow.AutocompleteComposerComponentsTaskID.Ref())
		if autocompleteResult != nil {
			if autocompleteResult.Error != "" {
				return autocompleteResult.Error, inspectionmetadata.Error, nil
			}
			if autocompleteResult.Hint != "" {
				return autocompleteResult.Hint, inspectionmetadata.Info, nil
			}
		}
		return "", inspectionmetadata.None, nil
	}).
	WithConverter(func(ctx context.Context, value []string) ([]string, error) {
		return value, nil
	}).
	Build()

Functions

func GenerateComposerLogsQuery

func GenerateComposerLogsQuery(projectID, location, environmentName string, selectedComponents []string) string

GenerateComposerLogsQuery generates a query for Composer environment logs.

func GenerateComposerLogsStructuredQuery

func GenerateComposerLogsStructuredQuery(projectID, location, environmentName string, selectedComponents []string) *logestimator.StructuredLogQuery

GenerateComposerLogsStructuredQuery generates a structured query for Composer environment logs.

func Register

func Register(registry coreinspection.InspectionTaskRegistry) error

Register registers all composer airflow inspection tasks to the registry.

Types

type DagProcessorState

type DagProcessorState struct {
	Reader *logutil.TabulateReader
}

DagProcessorState retains the parsing state using TabulateReader.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL