add_marketing_details_to_csv.mjs 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165
  1. #!/usr/bin/env node
  2. /* eslint-disable camelcase */
  3. import { scriptRunner } from './lib/ScriptRunner.mjs'
  4. import fs from 'node:fs'
  5. import minimist from 'minimist'
  6. import {
  7. db,
  8. ObjectId,
  9. READ_PREFERENCE_SECONDARY,
  10. } from '../app/src/infrastructure/mongodb.mjs'
  11. // https://github.com/import-js/eslint-plugin-import/issues/1810
  12. // eslint-disable-next-line import/no-unresolved
  13. import * as csv from 'csv/sync'
  14. function usage() {
  15. console.log(
  16. `
  17. This script enriches price and subscription data with user emails and first names, given a csv containing user IDs, outputs to /tmp/output.csv
  18. Usage:
  19. - Locally:
  20. docker compose exec web bash
  21. node scripts/add_marketing_details_to_csv.mjs [--input=<path>] [--output=<path>] [--batchSize=<number>]
  22. - On the server:
  23. rake run:pod[staging,web]
  24. node scripts/add_marketing_details_to_csv.mjs [--input=<path>] [--output=<path>] [--batchSize=<number>]
  25. exit
  26. kubectl cp web-standalone-prod-XXXXX:/tmp/output.csv ~/output.csv
  27. Options:
  28. --help Show this screen
  29. --input=<path> Input file path (default: input.csv)
  30. --output=<path> Output file path (default: /tmp/output.csv)
  31. --batchSize=<number> Number of users to be fetched in one query
  32. `
  33. )
  34. }
  35. function parseArgs() {
  36. const result = minimist(process.argv.slice(2), {
  37. string: ['input', 'output'],
  38. bool: ['help'],
  39. number: ['batchSize'],
  40. default: {
  41. help: false,
  42. input: 'input.csv',
  43. output: '/tmp/output.csv',
  44. batchSize: 1000,
  45. },
  46. })
  47. if (result.help) {
  48. usage()
  49. process.exit(0)
  50. }
  51. return result
  52. }
  53. async function processBatch(trackProgress, idBatch) {
  54. const result = {}
  55. try {
  56. const cursor = db.users.find(
  57. {
  58. _id: { $in: idBatch },
  59. },
  60. {
  61. projection: {
  62. _id: 1,
  63. first_name: 1,
  64. email: 1,
  65. },
  66. readPreference: READ_PREFERENCE_SECONDARY,
  67. }
  68. )
  69. for await (const doc of cursor) {
  70. result[doc._id.toString()] = {
  71. email: doc.email,
  72. first_name: doc.first_name,
  73. }
  74. }
  75. } catch (err) {
  76. await trackProgress(`ERROR Processing batch: ${err}`)
  77. }
  78. return result
  79. }
  80. async function enrichRecords(trackProgress, records, users) {
  81. const result = []
  82. for (const record of records) {
  83. const user = users[record.user_id]
  84. let email = ''
  85. let first_name = ''
  86. if (user) {
  87. if (user.email === '' || user.first_name === '') {
  88. await trackProgress(`WARNING Incomplete data for: ${record.user_id}`)
  89. }
  90. email = user.email
  91. first_name = user.first_name
  92. } else {
  93. await trackProgress(`WARNING Didn't find: ${record.user_id}`)
  94. }
  95. result.push({
  96. ...record,
  97. email,
  98. first_name,
  99. })
  100. }
  101. return result
  102. }
  103. async function main(trackProgress) {
  104. const args = parseArgs()
  105. const input = fs.readFileSync(args.input, 'utf8')
  106. const records = csv.parse(input, { columns: true, skipEmptyLines: true })
  107. await trackProgress(`INFO Starting to process ${records.length} records`)
  108. let objectIdBatch = []
  109. let recordBatch = []
  110. const outputRecords = []
  111. for (const record of records) {
  112. try {
  113. objectIdBatch.push(new ObjectId(record.user_id))
  114. recordBatch.push(record)
  115. } catch (e) {
  116. await trackProgress(`ERROR Skipping invalid user ID: ${record.user_id}`)
  117. outputRecords.push({ ...record, email: '', first_name: '' })
  118. }
  119. if (objectIdBatch.length >= args.batchSize) {
  120. const users = await processBatch(trackProgress, objectIdBatch)
  121. const enriched = await enrichRecords(trackProgress, recordBatch, users)
  122. outputRecords.push(...enriched)
  123. objectIdBatch = []
  124. recordBatch = []
  125. }
  126. }
  127. if (objectIdBatch.length > 0) {
  128. const users = await processBatch(trackProgress, objectIdBatch)
  129. const enriched = await enrichRecords(trackProgress, recordBatch, users)
  130. outputRecords.push(...enriched)
  131. }
  132. const output = csv.stringify(outputRecords, { header: true })
  133. fs.writeFileSync(args.output, output)
  134. await trackProgress(
  135. `INFO Finished processing of ${outputRecords.length} records`
  136. )
  137. }
  138. try {
  139. await scriptRunner(main)
  140. process.exit(0)
  141. } catch (error) {
  142. console.error(error)
  143. process.exit(1)
  144. }