Documentation
¶
Index ¶
- Variables
- func GenerateComposerLogsQuery(projectID, location, environmentName string, selectedComponents []string) string
- func GenerateComposerLogsStructuredQuery(projectID, location, environmentName string, selectedComponents []string) *logestimator.StructuredLogQuery
- func Register(registry coreinspection.InspectionTaskRegistry) error
- type DagProcessorState
Constants ¶
This section is empty.
Variables ¶
var AirflowDagProcessorManagerLogFilterTask = componentFilterTask(composerairflow.AirflowDagProcessorManagerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "dag-processor-manager")
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.
var AirflowDagProcessorManagerLogIngesterTask = inspectiontaskbase.NewGroupedLogIngesterTask( composerairflow.AirflowDagProcessorManagerLogIngesterTaskID, &dagProcessorManagerLogIngester{}, )
AirflowDagProcessorManagerLogIngesterTask is the task that ingests Airflow DAG processor manager logs.
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.
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" }, )
var AirflowOtherLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( composerairflow.AirflowOtherLogGrouperTaskID, composerairflow.AirflowOtherLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowOtherLogGrouperTask groups other Airflow logs.
var AirflowOtherLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( composerairflow.AirflowOtherLogIngesterTaskID, &otherLogIngester{}, )
AirflowOtherLogIngesterTask is the task that ingests other Airflow logs.
var AirflowOtherLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( composerairflow.AirflowOtherLogToTimelineMapperTaskID, &otherLogToTimelineMapper{}, )
AirflowOtherLogToTimelineMapperTask is the task that maps other Airflow logs to timeline events.
var AirflowSchedulerLogFilterTask = componentFilterTask(composerairflow.AirflowSchedulerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "airflow-scheduler")
var AirflowSchedulerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( composerairflow.AirflowSchedulerLogGrouperTaskID, composerairflow.AirflowSchedulerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowSchedulerLogGrouperTask groups Airflow scheduler logs.
var AirflowSchedulerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( composerairflow.AirflowSchedulerLogIngesterTaskID, &schedulerLogIngester{}, )
AirflowSchedulerLogIngesterTask is the task that ingests Airflow scheduler logs.
var AirflowSchedulerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( composerairflow.AirflowSchedulerLogToTimelineMapperTaskID, &schedulerLogToTimelineMapper{}, )
AirflowSchedulerLogToTimelineMapperTask is the task that maps Airflow scheduler logs to timeline events.
var AirflowWorkerLogFilterTask = componentFilterTask(composerairflow.AirflowWorkerLogFilterTaskID, composerairflow.ComposerLogsQueryTaskID.Ref(), "airflow-worker")
var AirflowWorkerLogGrouperTask = inspectiontaskbase.NewLogGrouperTask( composerairflow.AirflowWorkerLogGrouperTaskID, composerairflow.AirflowWorkerLogFilterTaskID.Ref(), func(ctx context.Context, l *log.Log) string { return "" }, )
AirflowWorkerLogGrouperTask groups Airflow worker logs.
var AirflowWorkerLogIngesterTask = inspectiontaskbase.NewLogIngesterTask( composerairflow.AirflowWorkerLogIngesterTaskID, &workerLogIngester{}, )
AirflowWorkerLogIngesterTask is the task that ingests Airflow worker logs.
var AirflowWorkerLogToTimelineMapperTask = inspectiontaskbase.NewLogToTimelineMapperTask( composerairflow.AirflowWorkerLogToTimelineMapperTaskID, &workerLogToTimelineMapper{}, )
AirflowWorkerLogToTimelineMapperTask is the task that maps Airflow worker logs to timeline events.
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 })
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.
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, ), )
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.