/
main.go
43 lines (40 loc) · 1.18 KB
/
main.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
package main
import (
"context"
"flag"
"github.com/fourstring/sheetfs/common_journal"
journal_example "github.com/fourstring/sheetfs/common_journal/example"
"github.com/fourstring/sheetfs/tests"
"log"
)
var num = flag.Int("n", 5, "Number of messages to produce")
var ckpt = flag.Bool("c", false, "whether to make a checkpoint after produce {num} messages")
var ctx = context.Background()
func main() {
flag.Parse()
writer, err := common_journal.NewWriter(journal_example.KafkaServer, journal_example.KafkaTopic)
if err != nil {
log.Fatal(err)
}
for i := 0; i < *num; i++ {
content := tests.RandStr(10)
err := writer.CommitEntry(ctx, []byte(content))
log.Printf("Primary: committed entry %s\n", content)
if err != nil {
log.Fatal(err)
}
}
if *ckpt {
writer.PrepareCheckpoint()
/*
Do real checkpoint operation here before call writer.Checkpoint()
*/
newStartOffset, err := writer.Checkpoint(ctx) // primary should persist newStartOffset, and it can
// start from here to recover data through journal entries.
if err != nil {
log.Fatal(err)
}
log.Printf("Primary: checkpoint successfully, and newStartOffset=%d\n", newStartOffset)
writer.ExitCheckpoint()
}
}