MCPcopy Create free account

hub / github.com/alibaba/MongoShake / functions

Functions933 in github.com/alibaba/MongoShake

↓ 6 callersMethodGTE
(v2 Version)
tools/pre-split/pre_split.go:85
↓ 6 callersFunctionNewDocumentReader
NewDocumentReader creates reader with mongodb url
collector/docsyncer/doc_reader.go:313
↓ 6 callersFunctionNewPacketV1
(packetType uint8, payload []byte)
tunnel/tcp_writer.go:77
↓ 6 callersFunctionRecordDuplicatedOplog
RecordDuplicatedOplog Write dup oplog in DB APPConflictDatabase
executor/duplicate.go:10
↓ 6 callersMethodSetFetchStage
(fetchStage int32)
collector/persister.go:72
↓ 6 callersFunctioncheckDefaultValue
()
cmd/collector/sanitize.go:85
↓ 6 callersMethoddoSync
use by full sync
collector/docsyncer/doc_executor.go:182
↓ 6 callersFunctiongetBsonType
(x interface{})
collector/filter/orphan_filter.go:219
↓ 6 callersFunctionprintCsOption
(ops *options.ChangeStreamOptions)
common/change_stream.go:159
↓ 6 callersFunctionreadOrder
(t *testing.T, order <-chan string)
cmd/collector/collector_test.go:143
↓ 5 callersFunctionInitNs
(specialNsList []string)
collector/filter/doc_filter.go:41
↓ 5 callersMethodLen
()
common/mix.go:24
↓ 5 callersMethodName
()
collector/coordinator/extra_job.go:70
↓ 5 callersFunctionPrometheusHandler
()
common/metric_prom.go:211
↓ 5 callersMethodPushToPendingQueue
(input []byte)
collector/persister.go:218
↓ 5 callersMethodSetReplStatus
(status uint64)
common/metric.go:338
↓ 5 callersFunctionTimeStampToInt64
(ts primitive.Timestamp)
common/mix.go:34
↓ 5 callersMethodUpdate
()
common/metric.go:35
↓ 5 callersMethodfilter
(log *oplog.PartialLog)
collector/filter/oplog_filter.go:381
↓ 5 callersFunctionformatArgs
(args ...any)
pkg/log/logger.go:165
↓ 5 callersFunctiongetTargetDelay
()
collector/batcher.go:40
↓ 5 callersFunctionlog_error
(message)
tools/mongodb-schema.py:26
↓ 5 callersFunctionlog_info
(message)
tools/mongodb-schema.py:22
↓ 5 callersMethodsetLastOplog
()
collector/batcher.go:394
↓ 5 callersMethodupdateCollectionProgressMetrics
()
collector/docsyncer/doc_syncer.go:649
↓ 5 callersMethodupdateLogsQueueMetric
(index int)
collector/syncer.go:570
↓ 4 callersFunctionCheckFcv
CheckFcv read the given file and parse the fcv do comparison
collector/configure/check.go:14
↓ 4 callersMethodCmp
(v2 Version)
tools/pre-split/pre_split.go:61
↓ 4 callersFunctionComputeHash
(data interface{})
collector/filter/orphan_filter.go:94
↓ 4 callersFunctionConvertBsonM2D
(input bson.M)
oplog/oplog.go:191
↓ 4 callersFunctionGetDbNamespace
GetDbNamespace return db namespace. return: @[]NS: namespace list, e.g., []{"a.b", "a.c"} @map[string][]string: db->collection map. e.g., "a"->[]strin
common/db_opertion.go:307
↓ 4 callersMethodIsAbort
IsAbort is true if the oplog entry had the abort command.
oplog/txn_meta.go:98
↓ 4 callersMethodIsCommit
IsCommit is true if the oplog entry was an abort command or was the final entry of an unprepared transaction.
oplog/txn_meta.go:109
↓ 4 callersMethodListen
()
common/http.go:36
↓ 4 callersMethodName
* * tunnel name */
tunnel/tunnel.go:150
↓ 4 callersFunctionNewBuffer
NewBuffer initializes a transaction oplog buffer.
oplog/txn_buffer.go:69
↓ 4 callersFunctionNewChangeStreamConn
(src string, mode string, fullDoc bool, specialDb string, filterFunc func(name string) bool, watchStartTi
common/change_stream.go:29
↓ 4 callersFunctionNewDocFilterList
()
collector/filter/doc_filter.go:138
↓ 4 callersMethodPurgeTxn
PurgeTxn closes any transaction streams in progress and deletes all oplog entries associated with a transaction. Must not be called concurrently with
oplog/txn_buffer.go:242
↓ 4 callersMethodSetLSN
(lsn int64)
common/metric.go:302
↓ 4 callersMethodString
()
collector/docsyncer/doc_reader.go:340
↓ 4 callersMethodSync
()
pkg/log/logger.go:280
↓ 4 callersFunctionWritePidById
(dir, id string)
common/mix.go:147
↓ 4 callersMethodaddIntoBatchGroup
addIntoBatchGroup isBarrier Barrier Oplogs(like DDL or Transaction) must execute sequentially and separately, send to batchGroup[0]
collector/batcher.go:413
↓ 4 callersMethodbyteSlice
(start, end int)
tools/mongo_id.go:37
↓ 4 callersFunctioncalculateSignature
(object interface{})
executor/collision_matrix.go:167
↓ 4 callersMethodcheckpoint
checkpoint calculate and update current checkpoint value. `flush` means whether force calculate & update checkpoint. if inputTs is given(> 0), use thi
collector/checkpoint.go:79
↓ 4 callersFunctioncrash
(msg string, errCode int)
cmd/receiver/receiver.go:130
↓ 4 callersFunctionfetchAllDocument
(conn *utils.MongoCommunityConn)
collector/docsyncer/doc_syncer_test.go:46
↓ 4 callersFunctionforwardCas
(v *int64, new int64)
common/metric.go:375
↓ 4 callersFunctiongenerate_conf
(tp, all, id, src_url, dst_url, dir_name, ckpt_time)
scripts/run_sys_test.py:38
↓ 4 callersFunctionintersection
(this bson.M, other bson.M)
executor/collision_matrix.go:300
↓ 4 callersMethodlog
(level zapcore.Level, message string)
pkg/log/logger.go:172
↓ 4 callersFunctionmarshalData
(input []bson.D)
collector/docsyncer/doc_syncer_test.go:32
↓ 4 callersFunctionmockTxnPartialOplogs
(startTs int64, normalOplog bool)
collector/batcher_test.go:142
↓ 4 callersMethodobservePutDelay
(ack int64, now time.Time)
collector/worker.go:277
↓ 4 callersMethodrelease
()
tunnel/tcp_writer.go:150
↓ 4 callersFunctionsendErrAndClose
sendErrAndClose is a utility for putting an error on a channel before closing.
oplog/txn_buffer.go:299
↓ 4 callersFunctionsocketTimeout
(socket *net.TCPConn, duration time.Duration)
tunnel/tcp_writer.go:253
↓ 4 callersMethodstartDeserializer
deserializer: fetch oplog from pending queue, parsed and then add into logs queue.
collector/syncer.go:442
↓ 4 callersMethodupdateBufferUsedMetric
()
collector/persister.go:231
↓ 4 callersMethodupdateJobsQueuedMetric
()
collector/worker.go:293
↓ 4 callersMethodupdateLSNLagMetrics
()
common/metric.go:352
↓ 4 callersMethodupdatePendingQueueMetric
(index int)
collector/syncer.go:558
↓ 4 callersMethodupdateUnackBufferMetric
()
collector/worker.go:305
↓ 3 callersMethodAddFilter
(incr uint64)
common/metric.go:264
↓ 3 callersMethodAddOp
Concurrency notes: We require that AddOp, GetTxnStream and PurgeTxn be called serially as part of orchestrating replay of oplog entries. The only me
oplog/txn_buffer.go:97
↓ 3 callersMethodAddSuccess
(incr uint64)
common/metric.go:239
↓ 3 callersMethodClose
()
common/change_stream.go:118
↓ 3 callersFunctionConvertBsonD2MExcept
(input bson.D, except map[string]struct{})
oplog/oplog.go:150
↓ 3 callersFunctionConvertEvent2Oplog
(input []byte, fullDoc bool)
oplog/change_stream_event.go:97
↓ 3 callersFunctionExtractInnerNs
ExtractInnerNs is used to extract the inner namespace of the txn PS: here only check the first inner op
oplog/txn_buffer.go:364
↓ 3 callersFunctionExtractInnerOps
ExtractInnerOps doc.applyOps[i].ts(Let ckpt use the last ts to judge complete) doc.applyOps[0 - n-1].ts = doc.ts - 1 doc.applyOps[n-1].ts = doc.ts
oplog/txn_buffer.go:321
↓ 3 callersFunctionExtractMongoTimestampCounter
(ts interface{})
common/mix.go:62
↓ 3 callersMethodGetFetchStage
()
collector/persister.go:77
↓ 3 callersFunctionGetIdOrNSFromOplog
(log *PartialLog)
oplog/hasher.go:122
↓ 3 callersFunctionGetNewestTimestampByUrl
(url string, fromMongoS bool, sslRootFile string)
common/db_opertion.go:138
↓ 3 callersMethodHasUniqueIndex
(queryCondition bson.M)
common/community_client.go:216
↓ 3 callersFunctionInitialLogger
InitialLogger initialize logger verbose: where log goes to: 0 - file,1 - file+stdout,2 - stdout
common/common.go:70
↓ 3 callersMethodIsFinal
IsFinal is true if the oplog entry is the closing entry of a transaction, i.e. if IsAbort or IsCommit is true.
oplog/txn_meta.go:115
↓ 3 callersMethodIsGood
()
common/community_client.go:178
↓ 3 callersFunctionIsMaster
()
quorum/quorum.go:53
↓ 3 callersFunctionIsSyncDataCommand
(operation string)
oplog/cmd_oplog.go:48
↓ 3 callersMethodIsTimeSeriesCollection
(dbName string, collName string)
common/community_client.go:268
↓ 3 callersMethodIterateFilter
(namespace string)
collector/filter/doc_filter.go:60
↓ 3 callersMethodName
()
collector/reader/reader.go:10
↓ 3 callersFunctionNewCollectionMetric
()
collector/docsyncer/metric.go:19
↓ 3 callersFunctionNewDBTransform
(transRule []string)
collector/transform/transform.go:55
↓ 3 callersFunctionNewGidFilter
(gids []string)
collector/filter/oplog_filter.go:149
↓ 3 callersFunctionNewMetric
var Metric *ReplicationMetric
common/metric.go:89
↓ 3 callersMethodPrepare
()
tunnel/kafka_writer.go:53
↓ 3 callersMethodRegister
()
common/sentinel.go:74
↓ 3 callersMethodResumeToken
()
common/change_stream.go:151
↓ 3 callersMethodSend
** * write the real tunnel message to tunnel. * * return the right ACK offset value with positive number. if AckRequired is set * this ACk off
tunnel/tunnel.go:140
↓ 3 callersMethodSend
(message *WMessage)
tunnel/kafka_writer.go:91
↓ 3 callersMethodSetLSNACK
(ack int64)
common/metric.go:308
↓ 3 callersMethodSetOplogDiskFinishTs
SetOplogDiskFinishTs OplogDiskQueueFinishTs and OplogDiskQueue won't take effect immediately, will be inserted in the next Update call.
collector/ckpt/ckpt_manager.go:111
↓ 3 callersMethodSetOplogDiskQueueName
(name string)
collector/ckpt/ckpt_manager.go:122
↓ 3 callersFunctionStartDropDestCollection
(nsSet map[utils.NS]struct{}, toConn *utils.MongoCommunityConn, nsTrans *transform.NamespaceTransform)
collector/docsyncer/doc_syncer.go:69
↓ 3 callersFunctionStartIndexSync
(indexMap map[utils.NS][]bson.D, toUrl string, nsTrans *transform.NamespaceTransform, background bool)
collector/docsyncer/doc_syncer.go:217
← previousnext →101–200 of 933, ranked by callers