readTableColumns reads table columns on applier
()
| 566 | |
| 567 | // readTableColumns reads table columns on applier |
| 568 | func (e *Extractor) readTableColumns() (err error) { |
| 569 | e.logger.Info("Examining table structure on extractor") |
| 570 | |
| 571 | // map parent -> child |
| 572 | fkParentMap := make(map[common.SchemaTable]map[common.SchemaTable]struct{}) |
| 573 | |
| 574 | for _, doDb := range e.replicateDoDb { |
| 575 | for _, tbCtx := range doDb.TableMap { |
| 576 | doTb := tbCtx.Table |
| 577 | tableColumns, fkParentTables, err := base.GetTableColumnsSqle(e.sqleContext, doTb.TableSchema, doTb.TableName) |
| 578 | if err != nil { |
| 579 | return err |
| 580 | } |
| 581 | doTb.OriginalTableColumns = tableColumns |
| 582 | doTb.ColumnMap, err = mysqlconfig.BuildColumnMapIndex(doTb.ColumnMapFrom, doTb.OriginalTableColumns.Ordinals) |
| 583 | if err != nil { |
| 584 | return err |
| 585 | } |
| 586 | |
| 587 | childST := common.SchemaTable{Schema: doTb.TableSchema, Table: doTb.TableName} |
| 588 | for _, fkpt := range fkParentTables { |
| 589 | schema := g.StringElse(fkpt.Schema.String(), doTb.TableSchema) |
| 590 | parentST := common.SchemaTable{Schema: schema, Table: fkpt.Name.String()} |
| 591 | if m, ok := fkParentMap[parentST]; ok { |
| 592 | m[childST] = struct{}{} |
| 593 | } else { |
| 594 | fkParentMap[parentST] = map[common.SchemaTable]struct{}{ |
| 595 | childST: {}, |
| 596 | } |
| 597 | } |
| 598 | } |
| 599 | } |
| 600 | } |
| 601 | |
| 602 | for _, db := range e.replicateDoDb { |
| 603 | for _, tbCtx := range db.TableMap { |
| 604 | tb := tbCtx.Table |
| 605 | st := common.SchemaTable{Schema: tb.TableSchema, Table: tb.TableName} |
| 606 | if m, ok := fkParentMap[st]; ok { |
| 607 | tbCtx.FKChildren = m |
| 608 | e.logger.Info("fk parent", "len", len(m), "schema", tb.TableSchema, "table", tb.TableName) |
| 609 | } |
| 610 | } |
| 611 | } |
| 612 | |
| 613 | return nil |
| 614 | } |
| 615 | |
| 616 | func (e *Extractor) initNatsPubClient(natsAddr string) (err error) { |
| 617 | e.logger.Debug("begin Connect nats server", "NatAddr", natsAddr) |
no test coverage detected