| package clickhouse |
|
|
| import ( |
| _ "embed" |
| "fmt" |
| "strings" |
| "time" |
|
|
| "github.com/huandu/go-sqlbuilder" |
|
|
| "github.com/openmeterio/openmeter/openmeter/streaming" |
| ) |
|
|
| |
| type createEventsTable struct { |
| Database string |
| EventsTableName string |
| } |
|
|
| func (d createEventsTable) toSQL() string { |
| tableName := getTableName(d.Database, d.EventsTableName) |
|
|
| sb := sqlbuilder.ClickHouse.NewCreateTableBuilder() |
| sb.CreateTable(tableName) |
| sb.IfNotExists() |
| sb.Define("namespace", "String") |
| sb.Define("id", "String") |
| sb.Define("type", "LowCardinality(String)") |
| sb.Define("subject", "String") |
| sb.Define("source", "String") |
| sb.Define("time", "DateTime") |
| sb.Define("data", "String") |
| sb.Define("ingested_at", "DateTime") |
| sb.Define("stored_at", "DateTime") |
| sb.Define(fmt.Sprintf("INDEX %s_stored_at stored_at TYPE minmax GRANULARITY 4", d.EventsTableName)) |
| sb.Define("store_row_id", "String") |
| sb.SQL("ENGINE = MergeTree") |
| sb.SQL("PARTITION BY toYYYYMM(time)") |
| |
| |
| |
| |
| |
| |
| sb.SQL("ORDER BY (namespace, type, subject, toStartOfHour(time))") |
|
|
| sql, _ := sb.Build() |
| return sql |
| } |
|
|
| |
| type queryEventsTable struct { |
| Database string |
| EventsTableName string |
| Namespace string |
| From time.Time |
| To *time.Time |
| IngestedAtFrom *time.Time |
| IngestedAtTo *time.Time |
| ID *string |
| Subject *string |
| Customers *[]streaming.Customer |
| Limit int |
| } |
|
|
| |
| |
| |
| func (d queryEventsTable) toCountRowSQL() (string, []interface{}) { |
| tableName := getTableName(d.Database, d.EventsTableName) |
|
|
| query := sqlbuilder.ClickHouse.NewSelectBuilder() |
| query.Select("count() as total") |
|
|
| query.From(tableName) |
|
|
| query.Where(query.Equal("namespace", d.Namespace)) |
| query.Where(query.GreaterEqualThan("time", d.From.Unix())) |
|
|
| if d.To != nil { |
| query.Where(query.LessThan("time", d.To.Unix())) |
| } |
|
|
| |
| var customers []streaming.Customer |
|
|
| if d.Customers != nil { |
| customers = *d.Customers |
| } |
|
|
| |
| var subjects []string |
|
|
| if d.Subject != nil { |
| subjects = append(subjects, *d.Subject) |
| } |
|
|
| query = subjectWhere(d.EventsTableName, subjects, query) |
| query = customersWhere(d.EventsTableName, customers, query) |
|
|
| sql, args := query.Build() |
| return sql, args |
| } |
|
|
| func (d queryEventsTable) toSQL() (string, []interface{}) { |
| tableName := getTableName(d.Database, d.EventsTableName) |
| query := sqlbuilder.ClickHouse.NewSelectBuilder() |
|
|
| |
| query.Select( |
| "id", |
| "type", |
| "subject", |
| "source", |
| "time", |
| "data", |
| "ingested_at", |
| "stored_at", |
| "store_row_id", |
| ) |
|
|
| |
| if d.Customers != nil { |
| query = selectCustomerIdColumn(d.EventsTableName, *d.Customers, query) |
| } |
|
|
| query.From(tableName) |
|
|
| |
| query.Where(query.Equal("namespace", d.Namespace)) |
| query.Where(query.GreaterEqualThan("time", d.From.Unix())) |
|
|
| if d.To != nil { |
| query.Where(query.LessThan("time", d.To.Unix())) |
| } |
| if d.IngestedAtFrom != nil { |
| query.Where(query.GreaterEqualThan("ingested_at", d.IngestedAtFrom.Unix())) |
| } |
| if d.IngestedAtTo != nil { |
| query.Where(query.LessThan("ingested_at", d.IngestedAtTo.Unix())) |
| } |
| if d.ID != nil { |
| query.Where(query.Like("id", fmt.Sprintf("%%%s%%", *d.ID))) |
| } |
|
|
| |
| var customers []streaming.Customer |
|
|
| if d.Customers != nil { |
| customers = *d.Customers |
| } |
|
|
| |
| var subjects []string |
|
|
| if d.Subject != nil { |
| subjects = append(subjects, *d.Subject) |
| } |
|
|
| query = subjectWhere(d.EventsTableName, subjects, query) |
| query = customersWhere(d.EventsTableName, customers, query) |
|
|
| |
| query.Desc().OrderBy("time") |
| query.Limit(d.Limit) |
|
|
| sql, args := query.Build() |
|
|
| return sql, args |
| } |
|
|
| type queryCountEvents struct { |
| Database string |
| EventsTableName string |
| Namespace string |
| From time.Time |
| } |
|
|
| func (d queryCountEvents) toSQL() (string, []interface{}) { |
| tableName := getTableName(d.Database, d.EventsTableName) |
|
|
| query := sqlbuilder.ClickHouse.NewSelectBuilder() |
| query.Select("count() as count", "subject") |
| query.From(tableName) |
|
|
| query.Where(query.Equal("namespace", d.Namespace)) |
| query.Where(query.GreaterEqualThan("time", d.From.Unix())) |
| query.GroupBy("subject") |
|
|
| sql, args := query.Build() |
| return sql, args |
| } |
|
|
| |
| type InsertEventsQuery struct { |
| Database string |
| EventsTableName string |
| Events []streaming.RawEvent |
| QuerySettings map[string]string |
| } |
|
|
| func (q InsertEventsQuery) ToSQL() (string, []interface{}) { |
| tableName := getTableName(q.Database, q.EventsTableName) |
|
|
| query := sqlbuilder.ClickHouse.NewInsertBuilder() |
| query.InsertInto(tableName) |
| query.Cols("namespace", "id", "type", "source", "subject", "time", "data", "ingested_at", "stored_at", "store_row_id") |
|
|
| |
| var settings []string |
| for key, value := range q.QuerySettings { |
| settings = append(settings, fmt.Sprintf("%s = %s", key, value)) |
| } |
|
|
| if len(settings) > 0 { |
| query.SQL(fmt.Sprintf("SETTINGS %s", strings.Join(settings, ", "))) |
| } |
|
|
| for _, event := range q.Events { |
| query.Values( |
| event.Namespace, |
| event.ID, |
| event.Type, |
| event.Source, |
| event.Subject, |
| event.Time, |
| event.Data, |
| event.IngestedAt, |
| event.StoredAt, |
| event.StoreRowID, |
| ) |
| } |
|
|
| sql, args := query.Build() |
| return sql, args |
| } |
|
|
| func getTableName(database string, tableName string) string { |
| return fmt.Sprintf("%s.%s", database, tableName) |
| } |
|
|