| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138 |
- package main
- import (
- "encoding/json"
- "flag"
- "fmt"
- "log"
- "os"
- amqp "github.com/rabbitmq/amqp091-go"
- "github.com/tealeg/xlsx"
- )
- func failOnError(err error, msg string) {
- if err != nil {
- log.Fatalf("%s: %v", msg, err)
- }
- }
- func main() {
- var (
- filePath = flag.String("file", "data.xlsx", "path to xlsx file")
- amqpURL = flag.String("amqp", "amqp://admin:ward-tootled-wick-intrados@localhost:5672/", "AMQP URL")
- queueName = flag.String("queue", "ingestion", "RabbitMQ queue name")
- sheetName = flag.String("sheet", "", "Sheet name (default: first sheet)")
- headerRow = flag.Bool("header", true, "treat first row as header keys")
- )
- flag.Parse()
- if _, err := os.Stat(*filePath); os.IsNotExist(err) {
- log.Fatalf("file not found: %s", *filePath)
- }
- xlFile, err := xlsx.OpenFile(*filePath)
- failOnError(err, "opening xlsx")
- var sh *xlsx.Sheet
- if *sheetName == "" {
- if len(xlFile.Sheets) == 0 {
- log.Fatalf("no sheets in workbook")
- }
- sh = xlFile.Sheets[0]
- } else {
- for _, s := range xlFile.Sheets {
- if s.Name == *sheetName {
- sh = s
- break
- }
- }
- if sh == nil {
- log.Fatalf("sheet %s not found", *sheetName)
- }
- }
- if len(sh.Rows) == 0 {
- log.Fatalf("sheet %s is empty", sh.Name)
- }
- // convert rows to [][]string
- rows := make([][]string, 0, len(sh.Rows))
- for _, r := range sh.Rows {
- cols := make([]string, 0, len(r.Cells))
- for _, c := range r.Cells {
- cols = append(cols, c.String())
- }
- rows = append(rows, cols)
- }
- conn, err := amqp.Dial(*amqpURL)
- failOnError(err, "connecting to amqp")
- defer conn.Close()
- ch, err := conn.Channel()
- failOnError(err, "opening channel")
- defer ch.Close()
- q, err := ch.QueueDeclare(
- *queueName,
- true,
- false,
- false,
- false,
- nil,
- )
- failOnError(err, "declaring queue")
- var headers []string
- start := 0
- if *headerRow {
- headers = rows[0]
- start = 1
- }
- sent := 0
- for i := start; i < len(rows); i++ {
- row := rows[i]
- var body []byte
- if *headerRow {
- m := make(map[string]string, len(headers))
- for j, h := range headers {
- val := ""
- if j < len(row) {
- val = row[j]
- }
- m[h] = val
- }
- body, err = json.Marshal(m)
- if err != nil {
- log.Printf("skipping row %d: marshal error: %v", i+1, err)
- continue
- }
- } else {
- body, err = json.Marshal(row)
- if err != nil {
- log.Printf("skipping row %d: marshal error: %v", i+1, err)
- continue
- }
- }
- err = ch.Publish(
- "", // exchange
- q.Name, // routing key = queue
- false,
- false,
- amqp.Publishing{
- ContentType: "application/json",
- Body: body,
- },
- )
- if err != nil {
- log.Printf("failed to publish row %d: %v", i+1, err)
- continue
- }
- sent++
- }
- fmt.Printf("sent %d messages from %s to queue %s\n", sent, *filePath, q.Name)
- }
|