Appearance
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)- 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.
- 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. - A second flow, triggered by the queue's own
queue.message.receivedevent, does the real work:file.parseturns the CSV into rows, a Loop card runsdb.writeonce 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 onA 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