Documentation
¶
Index ¶
- Constants
- Variables
- type AirflowTaskInstance
- func (a *AirflowTaskInstance) DagId() string
- func (a *AirflowTaskInstance) Host() string
- func (a *AirflowTaskInstance) MapIndex() string
- func (a *AirflowTaskInstance) ResourcePath() resourcepath.ResourcePath
- func (a *AirflowTaskInstance) RunId() string
- func (a *AirflowTaskInstance) Status() Tistate
- func (a *AirflowTaskInstance) TaskId() string
- func (a *AirflowTaskInstance) ToYaml() string
- type AirflowWorker
- type ComposerEnvironmentClusterFinder
- type ComposerEnvironmentIdentity
- type ComposerEnvironmentListFetcher
- type ComposerEnvironmentListFetcherImpl
- type ComposerFieldSet
- type ComposerFieldSetReader
- type ComposerTaskInstanceFieldSet
- type ComposerTaskInstanceFieldSetReader
- type ComposerWorkerTaskInstanceFieldSet
- type ComposerWorkerTaskInstanceFieldSetReader
- type EnvironmentClusterFinderImpl
- type Tistate
Constants ¶
const InspectionTypeID = "gcp-composer"
InspectionTypeID is the inspection type id for google cloud composer.
Variables ¶
var AirflowDagProcessorManagerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-dag-processor-manager")
AirflowDagProcessorManagerLogFilterTaskID is the task id for filtering Airflow DAG processor manager logs.
var AirflowDagProcessorManagerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-dag-processor-manager")
AirflowDagProcessorManagerLogGrouperTaskID is the task id for the task that groups Airflow DAG processor manager logs.
var AirflowDagProcessorManagerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-dag-processor-manager")
AirflowDagProcessorManagerLogIngesterTaskID is the task id for the task that ingests Airflow DAG processor manager logs.
var AirflowDagProcessorManagerLogSorterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "sorter-dag-processor-manager")
AirflowDagProcessorManagerLogSorterTaskID is the task id for the task that sorts Airflow DAG processor manager logs.
var AirflowDagProcessorManagerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-dag-processor-manager")
AirflowDagProcessorManagerLogToTimelineMapperTaskID is the task id for the task that maps Airflow DAG processor manager logs to timeline events.
var AirflowOtherLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-other")
AirflowOtherLogFilterTaskID is the task id for filtering other Airflow logs.
var AirflowOtherLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-other")
AirflowOtherLogGrouperTaskID is the task id for the task that groups other Airflow logs.
var AirflowOtherLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-other")
AirflowOtherLogIngesterTaskID is the task id for the task that ingests other Airflow logs.
var AirflowOtherLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-other")
AirflowOtherLogToTimelineMapperTaskID is the task id for the task that maps other Airflow logs to timeline events.
var AirflowSchedulerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-scheduler")
AirflowSchedulerLogFilterTaskID is the task id for filtering Airflow scheduler logs.
var AirflowSchedulerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-scheduler")
AirflowSchedulerLogGrouperTaskID is the task id for the task that groups Airflow scheduler logs.
var AirflowSchedulerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-scheduler")
AirflowSchedulerLogIngesterTaskID is the task id for the task that ingests Airflow scheduler logs.
var AirflowSchedulerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-scheduler")
AirflowSchedulerLogToTimelineMapperTaskID is the task id for the task that maps Airflow scheduler logs to timeline events.
var AirflowWorkerLogFilterTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "filter-worker")
AirflowWorkerLogFilterTaskID is the task id for filtering Airflow worker logs.
var AirflowWorkerLogGrouperTaskID = taskid.NewDefaultImplementationID[inspectiontaskbase.LogGroupMap](GoogleCloudComposerTaskIDPrefix + "grouper-worker")
AirflowWorkerLogGrouperTaskID is the task id for the task that groups Airflow worker logs.
var AirflowWorkerLogIngesterTaskID = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "ingester-worker")
AirflowWorkerLogIngesterTaskID is the task id for the task that ingests Airflow worker logs.
var AirflowWorkerLogToTimelineMapperTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "mapper-worker")
AirflowWorkerLogToTimelineMapperTaskID is the task id for the task that maps Airflow worker logs to timeline events.
var AutocompleteComposerClusterNamesTaskID = taskid.NewImplementationID(googlecloudk8scommon_contract.AutocompleteClusterIdentityTaskID.Ref(), "composer")
AutocompleteComposerClusterNamesTaskID is the task id for the task that autocompletes GKE cluster names created by Cloud Composer.
var AutocompleteComposerComponentsTaskID = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[string]](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-components")
AutocompleteComposerComponentsTaskID is the task id for autocompleting component names from Cloud Monitoring.
var AutocompleteComposerEnvironmentIdentityTaskID = taskid.NewDefaultImplementationID[*inspectioncore_contract.AutocompleteResult[ComposerEnvironmentIdentity]](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-environment-identities")
AutocompleteComposerEnvironmentIdentityTaskID is the task id for the task that autocompletes composer environment identities.
var AutocompleteComposerEnvironmentNamesTaskID taskid.TaskImplementationID[[]string] = taskid.NewDefaultImplementationID[[]string](GoogleCloudComposerTaskIDPrefix + "autocomplete/composer-environment-names")
AutocompleteComposerEnvironmentNamesTaskID is the task id for the task that autocompletes composer environment names.
var AutocompleteLocationForComposerEnvironmentTaskID = taskid.NewImplementationID(googlecloudcommon_contract.AutocompleteLocationTaskID.Ref(), "composer")
AutocompleteLocationForComposerEnvironmentTaskID is the task id for the task that autocompletes GKE cluster location from Composer environments.
var ClusterIdentityTaskID = taskid.NewDefaultImplementationID[googlecloudk8scommon_contract.GoogleCloudClusterIdentity](GoogleCloudComposerTaskIDPrefix + "cluster-identity")
ClusterIdentityTaskID is the task id for aliasing the cluster identity.
var ComposerClusterNamePrefixTaskID = taskid.NewImplementationID(googlecloudk8scommon_contract.ClusterNamePrefixTaskRef, "composer")
ComposerClusterNamePrefixTaskID is the task id for the task that returns the GKE cluster name prefix used by Cloud Composer.
var ComposerEnvironmentClusterFinderTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentClusterFinder](GoogleCloudComposerTaskIDPrefix + "composer-environment-cluster-finder")
ComposerEnvironmentClusterFinderTaskID is the task id for injecting ComposerEnvironmentClusterFinder instance.
var ComposerEnvironmentListFetcherTaskID = taskid.NewDefaultImplementationID[ComposerEnvironmentListFetcher](GoogleCloudComposerTaskIDPrefix + "composer-environment-list-fetcher")
ComposerEnvironmentListFetcherTaskID is the task id for injecting ComposerEnvironmentListFetcher instance.
var ComposerInspectionType = coreinspection.InspectionType{ Id: InspectionTypeID, Name: "Cloud Composer", Description: `Visualize logs related to Cloud Composer environment. Supports all GKE related logs(Cloud Composer v2) and Airflow logs(Airflow 2.0.0 or higher in any Cloud Composer version(v1-v2, partical v3))`, Icon: "assets/icons/composer.webp", Priority: math.MaxInt - 10, Labels: map[string]string{ inspectioncore_contract.InspectionTypeLabelKeyLogSource: "cloud_logging", inspectioncore_contract.InspectionTypeLabelKeyEnvironment: "googlecloud", inspectioncore_contract.InspectionTypeLabelKeyBasePlatform: "kubernetes", googlecloudcommon_contract.InspectionTypeLabelKeyClusterType: "gke", googlecloudcommon_contract.InspectionTypeLabelKeyProduct: "composer", }, }
ComposerInspectionType is the inspection type for google cloud composer.
var ComposerLogsFieldSetReadTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "fieldsetread")
ComposerLogsFieldSetReadTaskID is the task id for the task that reads fieldsets from composer logs.
var ComposerLogsQueryTaskID taskid.TaskImplementationID[[]*log.Log] = taskid.NewDefaultImplementationID[[]*log.Log](GoogleCloudComposerTaskIDPrefix + "query-composer-logs")
ComposerLogsQueryTaskID is the task id for the task that queries Logs from Cloud Logging.
var ComposerLogsTailTaskID = taskid.NewDefaultImplementationID[struct{}](GoogleCloudComposerTaskIDPrefix + "tail-composer-logs")
ComposerLogsTailTaskID is the task id for unifying composer logs feature.
var ErrEnvironmentClusterNotFound = errors.New("not found")
var GoogleCloudComposerTaskIDPrefix = "cloud.google.com/composer/"
GoogleCloudComposerTaskIDPrefix is the prefix for all task ids related to google cloud composer.
var InputComposerComponentsTaskID taskid.TaskImplementationID[[]string] = taskid.NewDefaultImplementationID[[]string](GoogleCloudComposerTaskIDPrefix + "input/composer/components")
InputComposerComponentsTaskID is the task id for selecting target Composer components.
var InputComposerEnvironmentNameTaskID taskid.TaskImplementationID[string] = taskid.NewDefaultImplementationID[string](GoogleCloudComposerTaskIDPrefix + "input/composer/environment_name")
InputComposerEnvironmentNameTaskID is the task id for the task that inputs composer environment name.
Functions ¶
This section is empty.
Types ¶
type AirflowTaskInstance ¶ added in v0.52.8
type AirflowTaskInstance struct {
// contains filtered or unexported fields
}
ref: https://github.com/apache/airflow/blob/main/airflow/models/taskinstance.py#L1187
func NewAirflowTaskInstance ¶ added in v0.52.8
func (*AirflowTaskInstance) DagId ¶ added in v0.52.8
func (a *AirflowTaskInstance) DagId() string
func (*AirflowTaskInstance) Host ¶ added in v0.52.8
func (a *AirflowTaskInstance) Host() string
func (*AirflowTaskInstance) MapIndex ¶ added in v0.52.8
func (a *AirflowTaskInstance) MapIndex() string
func (*AirflowTaskInstance) ResourcePath ¶ added in v0.52.8
func (a *AirflowTaskInstance) ResourcePath() resourcepath.ResourcePath
func (*AirflowTaskInstance) RunId ¶ added in v0.52.8
func (a *AirflowTaskInstance) RunId() string
func (*AirflowTaskInstance) Status ¶ added in v0.52.8
func (a *AirflowTaskInstance) Status() Tistate
func (*AirflowTaskInstance) TaskId ¶ added in v0.52.8
func (a *AirflowTaskInstance) TaskId() string
func (*AirflowTaskInstance) ToYaml ¶ added in v0.52.8
func (a *AirflowTaskInstance) ToYaml() string
type AirflowWorker ¶ added in v0.52.8
type AirflowWorker struct {
// contains filtered or unexported fields
}
func NewAirflowWorker ¶ added in v0.52.8
func NewAirflowWorker(host string) *AirflowWorker
func (*AirflowWorker) Host ¶ added in v0.52.8
func (a *AirflowWorker) Host() string
func (*AirflowWorker) ResourcePath ¶ added in v0.52.8
func (a *AirflowWorker) ResourcePath() resourcepath.ResourcePath
func (*AirflowWorker) ToYaml ¶ added in v0.52.8
func (a *AirflowWorker) ToYaml() string
type ComposerEnvironmentIdentity ¶ added in v0.52.0
ComposerEnvironmentIdentity represents the identity of a Cloud Composer environment.
type ComposerEnvironmentListFetcher ¶
type ComposerEnvironmentListFetcher interface {
GetEnvironmentNames(ctx context.Context, projectID, location string) ([]string, error)
}
ComposerEnvironmentListFetcher fetches the list of Cloud Composer environment names in project and the location.
type ComposerEnvironmentListFetcherImpl ¶
type ComposerEnvironmentListFetcherImpl struct{}
func (*ComposerEnvironmentListFetcherImpl) GetEnvironmentNames ¶
func (c *ComposerEnvironmentListFetcherImpl) GetEnvironmentNames(ctx context.Context, projectID string, location string) ([]string, error)
GetEnvironmentNames implements ComposerEnvironmentListFetcher.
type ComposerFieldSet ¶ added in v0.52.8
type ComposerFieldSet struct {
Component string // e.g. "worker", "scheduler", "dag-processor-manager"
WorkerID string
SchedulerID string
DagProcessorManagerID string
TriggererID string
WebserverID string
Subservice string
}
func (*ComposerFieldSet) Kind ¶ added in v0.52.8
func (c *ComposerFieldSet) Kind() string
type ComposerFieldSetReader ¶ added in v0.52.8
type ComposerFieldSetReader struct{}
func (*ComposerFieldSetReader) FieldSetKind ¶ added in v0.52.8
func (c *ComposerFieldSetReader) FieldSetKind() string
func (*ComposerFieldSetReader) Read ¶ added in v0.52.8
func (c *ComposerFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
type ComposerTaskInstanceFieldSet ¶ added in v0.52.8
type ComposerTaskInstanceFieldSet struct {
TaskInstance *AirflowTaskInstance
}
func (*ComposerTaskInstanceFieldSet) Kind ¶ added in v0.52.8
func (c *ComposerTaskInstanceFieldSet) Kind() string
type ComposerTaskInstanceFieldSetReader ¶ added in v0.52.8
type ComposerTaskInstanceFieldSetReader struct{}
func (*ComposerTaskInstanceFieldSetReader) FieldSetKind ¶ added in v0.52.8
func (c *ComposerTaskInstanceFieldSetReader) FieldSetKind() string
func (*ComposerTaskInstanceFieldSetReader) Read ¶ added in v0.52.8
func (c *ComposerTaskInstanceFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
type ComposerWorkerTaskInstanceFieldSet ¶ added in v0.52.8
type ComposerWorkerTaskInstanceFieldSet struct {
TaskInstance *AirflowTaskInstance
}
func (*ComposerWorkerTaskInstanceFieldSet) Kind ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSet) Kind() string
type ComposerWorkerTaskInstanceFieldSetReader ¶ added in v0.52.8
type ComposerWorkerTaskInstanceFieldSetReader struct{}
func (*ComposerWorkerTaskInstanceFieldSetReader) FieldSetKind ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSetReader) FieldSetKind() string
func (*ComposerWorkerTaskInstanceFieldSetReader) Read ¶ added in v0.52.8
func (c *ComposerWorkerTaskInstanceFieldSetReader) Read(reader *structured.NodeReader) (log.FieldSet, error)
type EnvironmentClusterFinderImpl ¶
type EnvironmentClusterFinderImpl struct{}
func (*EnvironmentClusterFinderImpl) GetGKEClusterName ¶
func (e *EnvironmentClusterFinderImpl) GetGKEClusterName(ctx context.Context, projectID string, environment string) (string, error)
GetGKEClusterName implements EnvironmentClusterFinder.
type Tistate ¶ added in v0.52.8
type Tistate string
const ( // ref: https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/tasks.html#task-instances TASKINSTANCE_NONE Tistate = "none" TASKINSTANCE_SCHEDULED Tistate = "scheduled" TASKINSTANCE_QUEUED Tistate = "queued" TASKINSTANCE_RUNNING Tistate = "running" TASKINSTANCE_SUCCESS Tistate = "success" TASKINSTANCE_SHUTDOWN Tistate = "shutdown" TASKINSTANCE_RESTARTING Tistate = "restarting" TASKINSTANCE_FAILED Tistate = "failed" TASKINSTANCE_SKIPPED Tistate = "skipped" TASKINSTANCE_UP_FOR_RETRY Tistate = "up_for_retry" TASKINSTANCE_DEFERRED Tistate = "deferred" TASKINSTANCE_UP_FOR_RESCHEDULE Tistate = "up_for_reschedule" TASKINSTANCE_REMOVED Tistate = "removed" TASKINSTANCE_UPSTREAM_FAILED Tistate = "upstream_failed" // Original States // // Zombie status for KHI view TASKINSTANCE_ZOMBIE Tistate = "zombie" )