MCPcopy Create free account
hub / github.com/crossjoin-io/crossjoin / createDataset

Method createDataset

api/datasets.go:143–198  ·  view source on GitHub ↗
(hash string, dataset config.Dataset)

Source from the content-addressed store, hash-verified

141}
142
143func (api *API) createDataset(hash string, dataset config.Dataset) error {
144 filename := filepath.Join(api.dataDir, dataset.ID+".db")
145
146 // Does the file exist? If so, remove it.
147 _, err := os.Stat(filename)
148 if err == nil {
149 log.Printf("`%s` already exists; removing", filename)
150 os.Remove(filename)
151 }
152
153 db, err := sql.Open("sqlite3", filename)
154 if err != nil {
155 return err
156 }
157 defer db.Close()
158 _, err = db.Exec("PRAGMA synchronous = OFF")
159 if err != nil {
160 return err
161 }
162 _, err = db.Exec("PRAGMA journal_mode = MEMORY")
163 if err != nil {
164 return err
165 }
166 _, err = db.Exec("PRAGMA cache_size = -2000000")
167 if err != nil {
168 return err
169 }
170
171 log.Printf("querying `%s`", dataset.DataSource.ID)
172 err = api.fetchSingle(hash, db, dataset.DataSource)
173 if err != nil {
174 return fmt.Errorf("fetch single: %w", err)
175 }
176
177 for _, join := range dataset.Joins {
178 log.Printf("querying `%s`", join.DataSource.ID)
179 err = api.fetchSingle(hash, db, join.DataSource)
180 if err != nil {
181 return fmt.Errorf("fetch single as part of join: %w", err)
182 }
183 }
184
185 joinClauses := ""
186 for _, join := range dataset.Joins {
187 joinColumns := []string{}
188 for _, cols := range join.Columns {
189 joinColumns = append(joinColumns, fmt.Sprintf(`%s."%s" = %s."%s"`, dataset.DataSource.ID, cols.LeftColumn, join.DataSource.ID, cols.RightColumn))
190 }
191 joinClauses += fmt.Sprintf(" %s %s ON %s", join.Type, join.DataSource.ID, strings.Join(joinColumns, " AND "))
192 }
193
194 log.Println("joining data")
195 joinQuery := fmt.Sprintf("CREATE TABLE %s AS SELECT * FROM %s %s", dataset.ID, dataset.DataSource.ID, joinClauses)
196 _, err = db.Exec(joinQuery)
197 return err
198}
199
200func (api *API) fetchSingle(hash string, dest *sql.DB, dataSource *config.DataSource) error {

Callers 1

refreshDatasetMethod · 0.95

Calls 1

fetchSingleMethod · 0.95

Tested by

no test coverage detected