Skip to content

A queue-driven file import that writes to a customer's database ​

When you need this ​

A partner drops a CSV on an SFTP server; each row needs writing into a customer-owned PostgreSQL database (not SPARK's own entities), and the work should not block the person who dropped the file, retry on its own if the database is briefly unreachable, and require a second person's sign-off before it runs for real. This combines four of the Integration Platform's own building blocks — see the full reference on the Integration and Connector Designer page.

The shape ​

folder watch (managed)  →  queue.publish  →  queue.message.received flow:
  file.parse (csv)  →  Loop over rows  →  db.write (bound parameters)
  1. A folder watch with a processed/failed folder takes each file that lands on the SFTP server and, once every flow listening for it has finished, archives it — see the ack/managed-file section of the same page.
  2. That first flow's only job is queue.publish — put the file's content on a queue and finish immediately, so the SFTP folder is never held up waiting for the database.
  3. A second flow, triggered by the queue's own queue.message.received event, does the real work: file.parse turns the CSV into rows, a Loop card runs db.write once per row with named parameters (never a string-built statement), and a failed row lands the whole message back on the queue to retry with a growing delay, then DEAD after the queue's attempts run out — visible on the Queues screen, with Replay.

The two flow definitions ​

intake — started by the folder watch's sftp.file.arrived event:

json
{
  "steps": [
    {
      "code": "enqueue",
      "type": "APP",
      "action": "queue.publish",
      "dependsOn": [],
      "params": { "queue": "orders-import", "message": { "path": "${input.path}", "connection": "${input.connection}" } }
    }
  ]
}

import-orders — started by queue.message.received:

json
{
  "steps": [
    {
      "code": "read",
      "type": "APP",
      "action": "sftp.read",
      "dependsOn": [],
      "params": { "connection": "${input.message.connection}", "path": "${input.message.path}" }
    },
    {
      "code": "rows",
      "type": "APP",
      "action": "file.parse",
      "dependsOn": ["read"],
      "params": { "content": "${steps.read.response.content}", "format": "csv" }
    },
    {
      "code": "writeRows",
      "type": "FOREACH",
      "forEach": "${steps.rows.response.rows}",
      "dependsOn": ["rows"],
      "steps": [
        {
          "code": "write",
          "type": "APP",
          "action": "db.write",
          "dependsOn": [],
          "params": {
            "connection": "orders-db",
            "sql": "INSERT INTO orders (order_no, customer, amount) VALUES (:orderNo, :customer, :amount)",
            "params": { "orderNo": "${item.order_no}", "customer": "${item.customer}", "amount": "${item.amount}" }
          }
        }
      ]
    }
  ]
}

Requiring a second person's sign-off before it runs for real ​

Both flows can be built and tested switched off. Before import-orders goes live against the real database:

bash
erp connection save orders-db --name "Orders database" --base-url "postgresql://reporting@db.example.com/orders" --type DATABASE --auth-type PASSWORD --secret-from-file ./orders-db-password.txt
erp flow require-golive-approval import-orders on --tenant 1
erp flow activate import-orders    # alice: this now makes a request instead of switching it on

A different person reviews and approves it — the person who asked can never also be the one who approves:

bash
erp api get "/api/v1/integration/golive-requests?status=PENDING" --tenant 1
erp api post /api/v1/integration/golive-requests/<id>/decide --tenant 1 --body '{"approved": true, "comment": "reviewed the statement and the retry settings"}'

Watching it live ​

bash
erp integration health                      # dead messages, failing watches, at a glance
erp queue messages orders-import --status DEAD
erp queue trace orders-import <messageId>   # every attempt's flow run, oldest first
erp queue replay orders-import <messageId>  # after fixing the real problem

Database connector, Queues, Connect to another system.