· 10 years ago · Sep 08, 2016, 04:38 AM
1package main
2
3import (
4 "context"
5 "log"
6 "os"
7 "sync"
8 "time"
9
10 "database/sql"
11 _ "github.com/lib/pq"
12)
13
14var (
15 db *sql.DB
16)
17
18type Other struct {
19 ID string
20 ThingID string
21 Value int
22
23 CreatedAt time.Time
24 UpdatedAt time.Time
25}
26
27type Thing struct {
28 ID string
29 Value int
30
31 Other Other
32
33 CreatedAt time.Time
34 UpdatedAt time.Time
35}
36
37type Retval struct {
38 Thing Thing
39
40 Value int
41 Err error
42}
43
44func init() {
45 var err error
46 if db, err = sql.Open("postgres", os.Getenv("DATABASE_URL")); nil != err {
47 log.Fatal(err)
48 }
49
50 sqls := []string{
51 `CREATE EXTENSION IF NOT EXISTS pgcrypto`,
52
53 `CREATE TABLE IF NOT EXISTS "things" (
54 "id" uuid DEFAULT gen_random_uuid() PRIMARY KEY,
55 "value" integer,
56 "created_at" timestamp without time zone NOT NULL,
57 "updated_at" timestamp without time zone NOT NULL
58 )`,
59
60 `CREATE TABLE IF NOT EXISTS "others" (
61 "id" uuid DEFAULT gen_random_uuid() PRIMARY KEY,
62
63 "thing_id" uuid NOT NULL UNIQUE REFERENCES things(id)
64 ON DELETE CASCADE
65 ON UPDATE RESTRICT,
66
67 "value" integer,
68
69 "created_at" timestamp without time zone NOT NULL,
70 "updated_at" timestamp without time zone NOT NULL
71 )`,
72 }
73
74 for _, sql := range sqls {
75 if _, err := db.Exec(sql); nil != err {
76 log.Fatal(err)
77 }
78 }
79}
80
81func Create(xs <-chan int) <-chan Retval {
82 group := &sync.WaitGroup{}
83
84 ch := make(chan Retval)
85
86 const sql1 = `INSERT INTO things (value, created_at, updated_at)
87 VALUES ($1, current_timestamp, current_timestamp)
88 RETURNING id`
89
90 const sql2 = `INSERT INTO others (thing_id, value, created_at, updated_at)
91 VALUES ($1, $2, current_timestamp, current_timestamp)`
92
93 output := func(x int) {
94 defer group.Done()
95
96 r := Retval{
97 Value: x,
98 }
99
100 tx, err := db.Begin()
101 if nil != err {
102 r.Err = err
103 ch <- r
104 return
105 }
106
107 var tableId string
108 if err = tx.QueryRow(sql1, x).Scan(&tableId); nil != err {
109 tx.Rollback()
110 r.Err = err
111 ch <- r
112 return
113 }
114
115 if _, err = tx.Exec(sql2, tableId, x); nil != err {
116 tx.Rollback()
117 r.Err = err
118 ch <- r
119 return
120 }
121
122 tx.Commit()
123 ch <- r
124 }
125
126 for x := range xs {
127 group.Add(1)
128 go output(x)
129 }
130
131 go func() {
132 group.Wait()
133 close(ch)
134 }()
135
136 return ch
137}
138
139func Values(ctx context.Context) <-chan int {
140 ch := make(chan int)
141
142 go func() {
143 defer close(ch)
144
145 for i := 1; i < 200; i++ {
146 select {
147 case ch <- i:
148 case <-ctx.Done():
149 return
150 }
151 }
152 }()
153
154 return ch
155}
156
157func Report() {
158 q := `SELECT others.thing_id, COUNT(others.thing_id)
159 FROM others
160 GROUP BY others.thing_id
161 HAVING 1 < COUNT(others.thing_id)`
162
163 rows, err := db.Query(q)
164 if nil != err {
165 log.Fatal(err)
166 }
167
168 defer rows.Close()
169
170 for rows.Next() {
171 var id string
172 var count int
173
174 rows.Scan(&id, &count)
175 log.Printf("[Report] Thing: %s, Count: %d\n", id, count)
176 }
177}
178
179func main() {
180 defer Report()
181
182 ctx, cancel := context.WithCancel(context.Background())
183 defer cancel()
184
185 values := Values(ctx)
186 ch := Create(values)
187
188 for r := range ch {
189 if nil != r.Err {
190 log.Printf("[Value: %4d] Error: %s\n", r.Value, r.Err.Error())
191 } else {
192 log.Printf("[Value: %4d] Success\n", r.Value)
193 }
194 }
195}