· 10 years ago · Aug 29, 2016, 05:30 PM
1<?xml version="1.0" ?>
2<template encoding-version="1.0">
3 <description>This template describes a flow where a CSV file (whose filename and content) contributes to the fields in a Cassandra table is processed, then CQL statements are constructed and executed.</description>
4 <groupId>d6aa94a0-0156-1000-71a4-a96b6da4672f</groupId>
5 <name>ConvertCSVtoCQL</name>
6 <snippet>
7 <connections>
8 <id>b7b54e60-d92f-4bb1-0000-000000000000</id>
9 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
10 <backPressureDataSizeThreshold>0 MB</backPressureDataSizeThreshold>
11 <backPressureObjectThreshold>0</backPressureObjectThreshold>
12 <destination>
13 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
14 <id>7969ae1c-8754-4f49-0000-000000000000</id>
15 <type>PROCESSOR</type>
16 </destination>
17 <flowFileExpiration>0 sec</flowFileExpiration>
18 <labelIndex>1</labelIndex>
19 <name></name>
20 <selectedRelationships>success</selectedRelationships>
21 <source>
22 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
23 <id>d794ca20-98c5-4fc5-0000-000000000000</id>
24 <type>PROCESSOR</type>
25 </source>
26 <zIndex>0</zIndex>
27 </connections>
28 <connections>
29 <id>d6ad6062-0156-1000-0000-000000000000</id>
30 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
31 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
32 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
33 <destination>
34 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
35 <id>d6ad3cb7-0156-1000-0000-000000000000</id>
36 <type>PROCESSOR</type>
37 </destination>
38 <flowFileExpiration>0 sec</flowFileExpiration>
39 <labelIndex>1</labelIndex>
40 <name></name>
41 <selectedRelationships>splits</selectedRelationships>
42 <source>
43 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
44 <id>d6abbdf3-0156-1000-0000-000000000000</id>
45 <type>PROCESSOR</type>
46 </source>
47 <zIndex>0</zIndex>
48 </connections>
49 <connections>
50 <id>d6aea457-0156-1000-0000-000000000000</id>
51 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
52 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
53 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
54 <destination>
55 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
56 <id>d6ae76e5-0156-1000-0000-000000000000</id>
57 <type>PROCESSOR</type>
58 </destination>
59 <flowFileExpiration>0 sec</flowFileExpiration>
60 <labelIndex>1</labelIndex>
61 <name></name>
62 <selectedRelationships>success</selectedRelationships>
63 <source>
64 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
65 <id>8e949039-0779-458d-0000-000000000000</id>
66 <type>PROCESSOR</type>
67 </source>
68 <zIndex>0</zIndex>
69 </connections>
70 <connections>
71 <id>d6aebdbc-0156-1000-0000-000000000000</id>
72 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
73 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
74 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
75 <destination>
76 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
77 <id>d6abbdf3-0156-1000-0000-000000000000</id>
78 <type>PROCESSOR</type>
79 </destination>
80 <flowFileExpiration>0 sec</flowFileExpiration>
81 <labelIndex>1</labelIndex>
82 <name></name>
83 <selectedRelationships>success</selectedRelationships>
84 <source>
85 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
86 <id>d6ae76e5-0156-1000-0000-000000000000</id>
87 <type>PROCESSOR</type>
88 </source>
89 <zIndex>0</zIndex>
90 </connections>
91 <connections>
92 <id>d6b4f5b7-0156-1000-0000-000000000000</id>
93 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
94 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
95 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
96 <destination>
97 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
98 <id>d6b4e56f-0156-1000-0000-000000000000</id>
99 <type>PROCESSOR</type>
100 </destination>
101 <flowFileExpiration>0 sec</flowFileExpiration>
102 <labelIndex>1</labelIndex>
103 <name></name>
104 <selectedRelationships>matched</selectedRelationships>
105 <source>
106 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
107 <id>d6ad3cb7-0156-1000-0000-000000000000</id>
108 <type>PROCESSOR</type>
109 </source>
110 <zIndex>0</zIndex>
111 </connections>
112 <connections>
113 <id>d6be6cc3-0156-1000-0000-000000000000</id>
114 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
115 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
116 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
117 <destination>
118 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
119 <id>d6be429d-0156-1000-0000-000000000000</id>
120 <type>PROCESSOR</type>
121 </destination>
122 <flowFileExpiration>0 sec</flowFileExpiration>
123 <labelIndex>1</labelIndex>
124 <name></name>
125 <selectedRelationships>success</selectedRelationships>
126 <source>
127 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
128 <id>d6b4e56f-0156-1000-0000-000000000000</id>
129 <type>PROCESSOR</type>
130 </source>
131 <zIndex>0</zIndex>
132 </connections>
133 <connections>
134 <id>d6d242d1-0156-1000-0000-000000000000</id>
135 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
136 <backPressureDataSizeThreshold>1 GB</backPressureDataSizeThreshold>
137 <backPressureObjectThreshold>10000</backPressureObjectThreshold>
138 <destination>
139 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
140 <id>7969ae1c-8754-4f49-0000-000000000000</id>
141 <type>PROCESSOR</type>
142 </destination>
143 <flowFileExpiration>0 sec</flowFileExpiration>
144 <labelIndex>1</labelIndex>
145 <name></name>
146 <selectedRelationships>success</selectedRelationships>
147 <source>
148 <groupId>d6aa94a0-0156-1000-0000-000000000000</groupId>
149 <id>d6be429d-0156-1000-0000-000000000000</id>
150 <type>PROCESSOR</type>
151 </source>
152 <zIndex>0</zIndex>
153 </connections>
154 <labels>
155 <id>d6e8d098-0156-1000-0000-000000000000</id>
156 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
157 <position>
158 <x>885.9999917116437</x>
159 <y>707.1400005203105</y>
160 </position>
161 <height>94.0</height>
162 <label>Creates the keyspace, drops the table,
163and creates the table. This is a helper
164for the example flow, should only be run
165once before the rest of the flow, and can
166be removed if the keyspace and table exist.
167<--</label>
168 <style>
169 <entry>
170 <key>font-size</key>
171 <value>12px</value>
172 </entry>
173 </style>
174 <width>246.99996948242188</width>
175 </labels>
176 <labels>
177 <id>d6e9804a-0156-1000-0000-000000000000</id>
178 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
179 <position>
180 <x>371.9999917116436</x>
181 <y>25.140004335007745</y>
182 </position>
183 <height>56.999996185302734</height>
184 <label>Generates a test CSV file with content.
185<--</label>
186 <style>
187 <entry>
188 <key>font-size</key>
189 <value>12px</value>
190 </entry>
191 </style>
192 <width>217.99996948242188</width>
193 </labels>
194 <labels>
195 <id>d6eadf91-0156-1000-0000-000000000000</id>
196 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
197 <position>
198 <x>368.9999917116436</x>
199 <y>228.14000052031048</y>
200 </position>
201 <height>67.0</height>
202 <label>The filename is station1_sensor2.csv,
203these segments are parsed to get
204fields for the CQL INSERTs.
205<--</label>
206 <style>
207 <entry>
208 <key>font-size</key>
209 <value>12px</value>
210 </entry>
211 </style>
212 <width>216.99996948242188</width>
213 </labels>
214 <labels>
215 <id>d6ec9645-0156-1000-0000-000000000000</id>
216 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
217 <position>
218 <x>1036.9999917116438</x>
219 <y>24.140004335007745</y>
220 </position>
221 <height>61.0</height>
222 <label>This attempts to strip the microseconds
223from the timestamp, but ends up truncating
224all fractional seconds.
225<--</label>
226 <style>
227 <entry>
228 <key>font-size</key>
229 <value>12px</value>
230 </entry>
231 </style>
232 <width>243.99996948242188</width>
233 </labels>
234 <labels>
235 <id>d6ed7b8f-0156-1000-0000-000000000000</id>
236 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
237 <position>
238 <x>1034.9999917116438</x>
239 <y>269.1400005203105</y>
240 </position>
241 <height>73.0</height>
242 <label>Builds a CQL INSERT statement with
243explicit values. Also could use UpdateAttribute
244in order to use prepared statements.
245<--</label>
246 <style>
247 <entry>
248 <key>font-size</key>
249 <value>12px</value>
250 </entry>
251 </style>
252 <width>271.9999694824219</width>
253 </labels>
254 <labels>
255 <id>d6ef37d2-0156-1000-0000-000000000000</id>
256 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
257 <position>
258 <x>1039.9999917116438</x>
259 <y>511.14000052031054</y>
260 </position>
261 <height>48.000003814697266</height>
262 <label>Executes the CQL statements
263(INSERT or DDL)
264<--</label>
265 <style>
266 <entry>
267 <key>font-size</key>
268 <value>12px</value>
269 </entry>
270 </style>
271 <width>174.99996948242188</width>
272 </labels>
273 <labels>
274 <id>d6f0d22f-0156-1000-0000-000000000000</id>
275 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
276 <position>
277 <x>365.9999917116436</x>
278 <y>619.1400005203105</y>
279 </position>
280 <height>48.000003814697266</height>
281 <label>Uses regular expressions to
282group the columns
283<--</label>
284 <style>
285 <entry>
286 <key>font-size</key>
287 <value>12px</value>
288 </entry>
289 </style>
290 <width>172.99996948242188</width>
291 </labels>
292 <processors>
293 <id>8e949039-0779-458d-0000-000000000000</id>
294 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
295 <position>
296 <x>0.0</x>
297 <y>1.0000114440917969</y>
298 </position>
299 <config>
300 <bulletinLevel>WARN</bulletinLevel>
301 <comments></comments>
302 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
303 <descriptors>
304 <entry>
305 <key>Script Engine</key>
306 <value>
307 <name>Script Engine</name>
308 </value>
309 </entry>
310 <entry>
311 <key>Script File</key>
312 <value>
313 <name>Script File</name>
314 </value>
315 </entry>
316 <entry>
317 <key>Script Body</key>
318 <value>
319 <name>Script Body</name>
320 </value>
321 </entry>
322 <entry>
323 <key>Module Directory</key>
324 <value>
325 <name>Module Directory</name>
326 </value>
327 </entry>
328 <entry>
329 <key>File Content</key>
330 <value>
331 <name>File Content</name>
332 </value>
333 </entry>
334 <entry>
335 <key>Evaluate Expressions in Content</key>
336 <value>
337 <name>Evaluate Expressions in Content</name>
338 </value>
339 </entry>
340 <entry>
341 <key>Filename</key>
342 <value>
343 <name>Filename</name>
344 </value>
345 </entry>
346 </descriptors>
347 <lossTolerant>false</lossTolerant>
348 <penaltyDuration>30 sec</penaltyDuration>
349 <properties>
350 <entry>
351 <key>Script Engine</key>
352 <value>Groovy</value>
353 </entry>
354 <entry>
355 <key>Script File</key>
356 </entry>
357 <entry>
358 <key>Script Body</key>
359 <value>class GenerateFlowFileWithContent implements Processor {
360
361 def REL_SUCCESS = new Relationship.Builder()
362 .name('success')
363 .description('The flow file with the specified content and/or filename was successfully transferred')
364 .build();
365
366 def CONTENT = new PropertyDescriptor.Builder()
367 .name('File Content').description('The content for the generated flow file')
368 .required(false).expressionLanguageSupported(true).addValidator(Validator.VALID).build()
369
370 def CONTENT_HAS_EL = new PropertyDescriptor.Builder()
371 .name('Evaluate Expressions in Content').description('Whether to evaluate NiFi Expression Language constructs within the content')
372 .required(true).allowableValues('true','false').defaultValue('false').build()
373
374 def FILENAME = new PropertyDescriptor.Builder()
375 .name('Filename').description('The name of the flow file to be stored in the filename attribute')
376 .required(false).expressionLanguageSupported(true).addValidator(StandardValidators.NON_EMPTY_VALIDATOR).build()
377
378 @Override
379 void initialize(ProcessorInitializationContext context) { }
380
381 @Override
382 Set<Relationship> getRelationships() { return [REL_SUCCESS] as Set }
383
384 @Override
385 void onTrigger(ProcessContext context, ProcessSessionFactory sessionFactory) throws ProcessException {
386 try {
387 def session = sessionFactory.createSession()
388 def flowFile = session.create()
389
390 def hasEL = context.getProperty(CONTENT_HAS_EL).asBoolean()
391 def contentProp = context.getProperty(CONTENT)
392 def content = (hasEL ? contentProp.evaluateAttributeExpressions().value : contentProp.value) ?: ''
393 def filename = context.getProperty(FILENAME)?.evaluateAttributeExpressions()?.getValue()
394
395 flowFile = session.write(flowFile, { outStream ->
396 outStream.write(content.getBytes("UTF-8"))
397 } as OutputStreamCallback)
398
399 if(filename != null) { flowFile = session.putAttribute(flowFile, 'filename', filename) }
400 // transfer
401 session.transfer(flowFile, REL_SUCCESS)
402 session.commit()
403 } catch(e) {
404 throw new ProcessException(e)
405 }
406 }
407
408 @Override
409 Collection<ValidationResult> validate(ValidationContext context) { return null }
410
411 @Override
412 PropertyDescriptor getPropertyDescriptor(String name) {
413 switch(name) {
414 case 'File Content': return CONTENT
415 case 'Evaluate Expressions in Content': return CONTENT_HAS_EL
416 case 'Filename': return FILENAME
417 default: return null
418 }
419 }
420
421 @Override
422 void onPropertyModified(PropertyDescriptor descriptor, String oldValue, String newValue) { }
423
424 @Override
425 List<PropertyDescriptor> getPropertyDescriptors() { return [CONTENT, CONTENT_HAS_EL, FILENAME] as List }
426
427 @Override
428 String getIdentifier() { return 'GenerateFlowFile-InvokeScriptedProcessor' }
429
430}
431
432processor = new GenerateFlowFileWithContent()</value>
433 </entry>
434 <entry>
435 <key>Module Directory</key>
436 </entry>
437 <entry>
438 <key>File Content</key>
439 <value>2016-05-04 03:02:01.001000+0000;0;
4402016-05-04 03:02:01.002000+0000;0.1234;
4412016-05-04 03:02:01.003000+0000;0.2345;</value>
442 </entry>
443 <entry>
444 <key>Evaluate Expressions in Content</key>
445 <value>false</value>
446 </entry>
447 <entry>
448 <key>Filename</key>
449 <value>station1_sensor2.csv</value>
450 </entry>
451 </properties>
452 <runDurationMillis>0</runDurationMillis>
453 <schedulingPeriod>30 sec</schedulingPeriod>
454 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
455 <yieldDuration>1 sec</yieldDuration>
456 </config>
457 <name>GenerateFlowFileWithContent</name>
458 <relationships>
459 <autoTerminate>false</autoTerminate>
460 <name>success</name>
461 </relationships>
462 <style></style>
463 <type>org.apache.nifi.processors.script.InvokeScriptedProcessor</type>
464 </processors>
465 <processors>
466 <id>d6abbdf3-0156-1000-0000-000000000000</id>
467 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
468 <position>
469 <x>4.0</x>
470 <y>391.99999618530273</y>
471 </position>
472 <config>
473 <bulletinLevel>WARN</bulletinLevel>
474 <comments></comments>
475 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
476 <descriptors>
477 <entry>
478 <key>Line Split Count</key>
479 <value>
480 <name>Line Split Count</name>
481 </value>
482 </entry>
483 <entry>
484 <key>Maximum Fragment Size</key>
485 <value>
486 <name>Maximum Fragment Size</name>
487 </value>
488 </entry>
489 <entry>
490 <key>Header Line Count</key>
491 <value>
492 <name>Header Line Count</name>
493 </value>
494 </entry>
495 <entry>
496 <key>Header Line Marker Characters</key>
497 <value>
498 <name>Header Line Marker Characters</name>
499 </value>
500 </entry>
501 <entry>
502 <key>Remove Trailing Newlines</key>
503 <value>
504 <name>Remove Trailing Newlines</name>
505 </value>
506 </entry>
507 </descriptors>
508 <lossTolerant>false</lossTolerant>
509 <penaltyDuration>30 sec</penaltyDuration>
510 <properties>
511 <entry>
512 <key>Line Split Count</key>
513 <value>1</value>
514 </entry>
515 <entry>
516 <key>Maximum Fragment Size</key>
517 </entry>
518 <entry>
519 <key>Header Line Count</key>
520 <value>0</value>
521 </entry>
522 <entry>
523 <key>Header Line Marker Characters</key>
524 </entry>
525 <entry>
526 <key>Remove Trailing Newlines</key>
527 <value>true</value>
528 </entry>
529 </properties>
530 <runDurationMillis>0</runDurationMillis>
531 <schedulingPeriod>0 sec</schedulingPeriod>
532 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
533 <yieldDuration>1 sec</yieldDuration>
534 </config>
535 <name>Split CSV Lines</name>
536 <relationships>
537 <autoTerminate>true</autoTerminate>
538 <name>failure</name>
539 </relationships>
540 <relationships>
541 <autoTerminate>true</autoTerminate>
542 <name>original</name>
543 </relationships>
544 <relationships>
545 <autoTerminate>false</autoTerminate>
546 <name>splits</name>
547 </relationships>
548 <style></style>
549 <type>org.apache.nifi.processors.standard.SplitText</type>
550 </processors>
551 <processors>
552 <id>d6ad3cb7-0156-1000-0000-000000000000</id>
553 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
554 <position>
555 <x>1.0</x>
556 <y>612.9999961853027</y>
557 </position>
558 <config>
559 <bulletinLevel>WARN</bulletinLevel>
560 <comments></comments>
561 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
562 <descriptors>
563 <entry>
564 <key>Character Set</key>
565 <value>
566 <name>Character Set</name>
567 </value>
568 </entry>
569 <entry>
570 <key>Maximum Buffer Size</key>
571 <value>
572 <name>Maximum Buffer Size</name>
573 </value>
574 </entry>
575 <entry>
576 <key>Maximum Capture Group Length</key>
577 <value>
578 <name>Maximum Capture Group Length</name>
579 </value>
580 </entry>
581 <entry>
582 <key>Enable Canonical Equivalence</key>
583 <value>
584 <name>Enable Canonical Equivalence</name>
585 </value>
586 </entry>
587 <entry>
588 <key>Enable Case-insensitive Matching</key>
589 <value>
590 <name>Enable Case-insensitive Matching</name>
591 </value>
592 </entry>
593 <entry>
594 <key>Permit Whitespace and Comments in Pattern</key>
595 <value>
596 <name>Permit Whitespace and Comments in Pattern</name>
597 </value>
598 </entry>
599 <entry>
600 <key>Enable DOTALL Mode</key>
601 <value>
602 <name>Enable DOTALL Mode</name>
603 </value>
604 </entry>
605 <entry>
606 <key>Enable Literal Parsing of the Pattern</key>
607 <value>
608 <name>Enable Literal Parsing of the Pattern</name>
609 </value>
610 </entry>
611 <entry>
612 <key>Enable Multiline Mode</key>
613 <value>
614 <name>Enable Multiline Mode</name>
615 </value>
616 </entry>
617 <entry>
618 <key>Enable Unicode-aware Case Folding</key>
619 <value>
620 <name>Enable Unicode-aware Case Folding</name>
621 </value>
622 </entry>
623 <entry>
624 <key>Enable Unicode Predefined Character Classes</key>
625 <value>
626 <name>Enable Unicode Predefined Character Classes</name>
627 </value>
628 </entry>
629 <entry>
630 <key>Enable Unix Lines Mode</key>
631 <value>
632 <name>Enable Unix Lines Mode</name>
633 </value>
634 </entry>
635 <entry>
636 <key>Include Capture Group 0</key>
637 <value>
638 <name>Include Capture Group 0</name>
639 </value>
640 </entry>
641 <entry>
642 <key>column</key>
643 <value>
644 <name>column</name>
645 </value>
646 </entry>
647 </descriptors>
648 <lossTolerant>false</lossTolerant>
649 <penaltyDuration>30 sec</penaltyDuration>
650 <properties>
651 <entry>
652 <key>Character Set</key>
653 <value>UTF-8</value>
654 </entry>
655 <entry>
656 <key>Maximum Buffer Size</key>
657 <value>1 MB</value>
658 </entry>
659 <entry>
660 <key>Maximum Capture Group Length</key>
661 <value>1024</value>
662 </entry>
663 <entry>
664 <key>Enable Canonical Equivalence</key>
665 <value>false</value>
666 </entry>
667 <entry>
668 <key>Enable Case-insensitive Matching</key>
669 <value>false</value>
670 </entry>
671 <entry>
672 <key>Permit Whitespace and Comments in Pattern</key>
673 <value>false</value>
674 </entry>
675 <entry>
676 <key>Enable DOTALL Mode</key>
677 <value>false</value>
678 </entry>
679 <entry>
680 <key>Enable Literal Parsing of the Pattern</key>
681 <value>false</value>
682 </entry>
683 <entry>
684 <key>Enable Multiline Mode</key>
685 <value>false</value>
686 </entry>
687 <entry>
688 <key>Enable Unicode-aware Case Folding</key>
689 <value>false</value>
690 </entry>
691 <entry>
692 <key>Enable Unicode Predefined Character Classes</key>
693 <value>false</value>
694 </entry>
695 <entry>
696 <key>Enable Unix Lines Mode</key>
697 <value>false</value>
698 </entry>
699 <entry>
700 <key>Include Capture Group 0</key>
701 <value>true</value>
702 </entry>
703 <entry>
704 <key>column</key>
705 <value>([^;]*?);([^;]*?);([^;]*?)</value>
706 </entry>
707 </properties>
708 <runDurationMillis>0</runDurationMillis>
709 <schedulingPeriod>0 sec</schedulingPeriod>
710 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
711 <yieldDuration>1 sec</yieldDuration>
712 </config>
713 <name>Extract Column Values</name>
714 <relationships>
715 <autoTerminate>false</autoTerminate>
716 <name>matched</name>
717 </relationships>
718 <relationships>
719 <autoTerminate>true</autoTerminate>
720 <name>unmatched</name>
721 </relationships>
722 <style></style>
723 <type>org.apache.nifi.processors.standard.ExtractText</type>
724 </processors>
725 <processors>
726 <id>d6ae76e5-0156-1000-0000-000000000000</id>
727 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
728 <position>
729 <x>4.0</x>
730 <y>203.0000114440918</y>
731 </position>
732 <config>
733 <bulletinLevel>WARN</bulletinLevel>
734 <comments></comments>
735 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
736 <descriptors>
737 <entry>
738 <key>Delete Attributes Expression</key>
739 <value>
740 <name>Delete Attributes Expression</name>
741 </value>
742 </entry>
743 <entry>
744 <key>sensor.name</key>
745 <value>
746 <name>sensor.name</name>
747 </value>
748 </entry>
749 <entry>
750 <key>station.name</key>
751 <value>
752 <name>station.name</name>
753 </value>
754 </entry>
755 </descriptors>
756 <lossTolerant>false</lossTolerant>
757 <penaltyDuration>30 sec</penaltyDuration>
758 <properties>
759 <entry>
760 <key>Delete Attributes Expression</key>
761 </entry>
762 <entry>
763 <key>sensor.name</key>
764 <value>${filename:substringAfter('_'):substringBefore('.csv')}</value>
765 </entry>
766 <entry>
767 <key>station.name</key>
768 <value>${filename:substringBefore('_')}</value>
769 </entry>
770 </properties>
771 <runDurationMillis>0</runDurationMillis>
772 <schedulingPeriod>0 sec</schedulingPeriod>
773 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
774 <yieldDuration>1 sec</yieldDuration>
775 </config>
776 <name>Parse Filename</name>
777 <relationships>
778 <autoTerminate>false</autoTerminate>
779 <name>success</name>
780 </relationships>
781 <style></style>
782 <type>org.apache.nifi.processors.attributes.UpdateAttribute</type>
783 </processors>
784 <processors>
785 <id>d6b4e56f-0156-1000-0000-000000000000</id>
786 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
787 <position>
788 <x>656.0</x>
789 <y>0.0</y>
790 </position>
791 <config>
792 <bulletinLevel>WARN</bulletinLevel>
793 <comments></comments>
794 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
795 <descriptors>
796 <entry>
797 <key>Delete Attributes Expression</key>
798 <value>
799 <name>Delete Attributes Expression</name>
800 </value>
801 </entry>
802 <entry>
803 <key>tps</key>
804 <value>
805 <name>tps</name>
806 </value>
807 </entry>
808 </descriptors>
809 <lossTolerant>false</lossTolerant>
810 <penaltyDuration>30 sec</penaltyDuration>
811 <properties>
812 <entry>
813 <key>Delete Attributes Expression</key>
814 </entry>
815 <entry>
816 <key>tps</key>
817 <value>${column.1:toDate('yyyy-MM-dd HH:mm:ss.SSSSSSZ'):format('yyyy-MM-dd HH:mm:ss.SSSZ')}</value>
818 </entry>
819 </properties>
820 <runDurationMillis>0</runDurationMillis>
821 <schedulingPeriod>0 sec</schedulingPeriod>
822 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
823 <yieldDuration>1 sec</yieldDuration>
824 </config>
825 <name>Convert Timestamp</name>
826 <relationships>
827 <autoTerminate>false</autoTerminate>
828 <name>success</name>
829 </relationships>
830 <style></style>
831 <type>org.apache.nifi.processors.attributes.UpdateAttribute</type>
832 </processors>
833 <processors>
834 <id>d6be429d-0156-1000-0000-000000000000</id>
835 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
836 <position>
837 <x>652.9999999999999</x>
838 <y>245.99999618530273</y>
839 </position>
840 <config>
841 <bulletinLevel>WARN</bulletinLevel>
842 <comments></comments>
843 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
844 <descriptors>
845 <entry>
846 <key>Regular Expression</key>
847 <value>
848 <name>Regular Expression</name>
849 </value>
850 </entry>
851 <entry>
852 <key>Replacement Value</key>
853 <value>
854 <name>Replacement Value</name>
855 </value>
856 </entry>
857 <entry>
858 <key>Character Set</key>
859 <value>
860 <name>Character Set</name>
861 </value>
862 </entry>
863 <entry>
864 <key>Maximum Buffer Size</key>
865 <value>
866 <name>Maximum Buffer Size</name>
867 </value>
868 </entry>
869 <entry>
870 <key>Replacement Strategy</key>
871 <value>
872 <name>Replacement Strategy</name>
873 </value>
874 </entry>
875 <entry>
876 <key>Evaluation Mode</key>
877 <value>
878 <name>Evaluation Mode</name>
879 </value>
880 </entry>
881 </descriptors>
882 <lossTolerant>false</lossTolerant>
883 <penaltyDuration>30 sec</penaltyDuration>
884 <properties>
885 <entry>
886 <key>Regular Expression</key>
887 <value>(?s:^.*$)</value>
888 </entry>
889 <entry>
890 <key>Replacement Value</key>
891 <value>insert into sensor_data (station_id, sensor_id, tps, val) values ('${station.name}', '${sensor.name}', '${tps}', ${column.2})</value>
892 </entry>
893 <entry>
894 <key>Character Set</key>
895 <value>UTF-8</value>
896 </entry>
897 <entry>
898 <key>Maximum Buffer Size</key>
899 <value>1 MB</value>
900 </entry>
901 <entry>
902 <key>Replacement Strategy</key>
903 <value>Regex Replace</value>
904 </entry>
905 <entry>
906 <key>Evaluation Mode</key>
907 <value>Entire text</value>
908 </entry>
909 </properties>
910 <runDurationMillis>0</runDurationMillis>
911 <schedulingPeriod>0 sec</schedulingPeriod>
912 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
913 <yieldDuration>1 sec</yieldDuration>
914 </config>
915 <name>ReplaceText</name>
916 <relationships>
917 <autoTerminate>true</autoTerminate>
918 <name>failure</name>
919 </relationships>
920 <relationships>
921 <autoTerminate>false</autoTerminate>
922 <name>success</name>
923 </relationships>
924 <style></style>
925 <type>org.apache.nifi.processors.standard.ReplaceText</type>
926 </processors>
927 <processors>
928 <id>d794ca20-98c5-4fc5-0000-000000000000</id>
929 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
930 <position>
931 <x>503.4999999999999</x>
932 <y>687.7600114440918</y>
933 </position>
934 <config>
935 <bulletinLevel>WARN</bulletinLevel>
936 <comments></comments>
937 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
938 <descriptors>
939 <entry>
940 <key>Script Engine</key>
941 <value>
942 <name>Script Engine</name>
943 </value>
944 </entry>
945 <entry>
946 <key>Script File</key>
947 <value>
948 <name>Script File</name>
949 </value>
950 </entry>
951 <entry>
952 <key>Script Body</key>
953 <value>
954 <name>Script Body</name>
955 </value>
956 </entry>
957 <entry>
958 <key>Module Directory</key>
959 <value>
960 <name>Module Directory</name>
961 </value>
962 </entry>
963 </descriptors>
964 <lossTolerant>false</lossTolerant>
965 <penaltyDuration>30 sec</penaltyDuration>
966 <properties>
967 <entry>
968 <key>Script Engine</key>
969 <value>Groovy</value>
970 </entry>
971 <entry>
972 <key>Script File</key>
973 </entry>
974 <entry>
975 <key>Script Body</key>
976 <value>DDL =
977"""CREATE KEYSPACE IF NOT EXISTS data with replication = {'class' : 'SimpleStrategy', 'replication_factor' : 3};
978 DROP TABLE IF EXISTS sensor_data;
979 CREATE TABLE sensor_data (station_id text, sensor_id text, tps timestamp, val float, PRIMARY KEY ((station_id, sensor_id), tps));
980"""
981
982DDL.eachLine { ddl ->
983 flowFile = session.create()
984 flowFile = session.write(flowFile, { outStream ->
985 outStream.write(ddl.getBytes('UTF-8'))
986 } as OutputStreamCallback)
987 session.transfer(flowFile, REL_SUCCESS)
988}</value>
989 </entry>
990 <entry>
991 <key>Module Directory</key>
992 </entry>
993 </properties>
994 <runDurationMillis>0</runDurationMillis>
995 <schedulingPeriod>60 sec</schedulingPeriod>
996 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
997 <yieldDuration>1 sec</yieldDuration>
998 </config>
999 <name>ExecuteScript</name>
1000 <relationships>
1001 <autoTerminate>true</autoTerminate>
1002 <name>failure</name>
1003 </relationships>
1004 <relationships>
1005 <autoTerminate>false</autoTerminate>
1006 <name>success</name>
1007 </relationships>
1008 <style>
1009 <entry>
1010 <key>background-color</key>
1011 <value>#000000</value>
1012 </entry>
1013 </style>
1014 <type>org.apache.nifi.processors.script.ExecuteScript</type>
1015 </processors>
1016 <processors>
1017 <id>7969ae1c-8754-4f49-0000-000000000000</id>
1018 <parentGroupId>d6aa94a0-0156-1000-0000-000000000000</parentGroupId>
1019 <position>
1020 <x>654.9999999999999</x>
1021 <y>472.73999618530274</y>
1022 </position>
1023 <config>
1024 <bulletinLevel>WARN</bulletinLevel>
1025 <comments></comments>
1026 <concurrentlySchedulableTaskCount>1</concurrentlySchedulableTaskCount>
1027 <descriptors>
1028 <entry>
1029 <key>Cassandra Contact Points</key>
1030 <value>
1031 <name>Cassandra Contact Points</name>
1032 </value>
1033 </entry>
1034 <entry>
1035 <key>Keyspace</key>
1036 <value>
1037 <name>Keyspace</name>
1038 </value>
1039 </entry>
1040 <entry>
1041 <key>SSL Context Service</key>
1042 <value>
1043 <identifiesControllerService>org.apache.nifi.ssl.SSLContextService</identifiesControllerService>
1044 <name>SSL Context Service</name>
1045 </value>
1046 </entry>
1047 <entry>
1048 <key>Client Auth</key>
1049 <value>
1050 <name>Client Auth</name>
1051 </value>
1052 </entry>
1053 <entry>
1054 <key>Username</key>
1055 <value>
1056 <name>Username</name>
1057 </value>
1058 </entry>
1059 <entry>
1060 <key>Password</key>
1061 <value>
1062 <name>Password</name>
1063 </value>
1064 </entry>
1065 <entry>
1066 <key>Consistency Level</key>
1067 <value>
1068 <name>Consistency Level</name>
1069 </value>
1070 </entry>
1071 <entry>
1072 <key>Character Set</key>
1073 <value>
1074 <name>Character Set</name>
1075 </value>
1076 </entry>
1077 <entry>
1078 <key>Max Wait Time</key>
1079 <value>
1080 <name>Max Wait Time</name>
1081 </value>
1082 </entry>
1083 </descriptors>
1084 <lossTolerant>false</lossTolerant>
1085 <penaltyDuration>30 sec</penaltyDuration>
1086 <properties>
1087 <entry>
1088 <key>Cassandra Contact Points</key>
1089 <value>127.0.0.1:9042</value>
1090 </entry>
1091 <entry>
1092 <key>Keyspace</key>
1093 <value>data</value>
1094 </entry>
1095 <entry>
1096 <key>SSL Context Service</key>
1097 </entry>
1098 <entry>
1099 <key>Client Auth</key>
1100 <value>REQUIRED</value>
1101 </entry>
1102 <entry>
1103 <key>Username</key>
1104 <value>cassandra</value>
1105 </entry>
1106 <entry>
1107 <key>Password</key>
1108 </entry>
1109 <entry>
1110 <key>Consistency Level</key>
1111 <value>ONE</value>
1112 </entry>
1113 <entry>
1114 <key>Character Set</key>
1115 <value>UTF-8</value>
1116 </entry>
1117 <entry>
1118 <key>Max Wait Time</key>
1119 <value>0 seconds</value>
1120 </entry>
1121 </properties>
1122 <runDurationMillis>0</runDurationMillis>
1123 <schedulingPeriod>0 sec</schedulingPeriod>
1124 <schedulingStrategy>TIMER_DRIVEN</schedulingStrategy>
1125 <yieldDuration>1 sec</yieldDuration>
1126 </config>
1127 <name>PutCassandraQL</name>
1128 <relationships>
1129 <autoTerminate>true</autoTerminate>
1130 <name>failure</name>
1131 </relationships>
1132 <relationships>
1133 <autoTerminate>true</autoTerminate>
1134 <name>retry</name>
1135 </relationships>
1136 <relationships>
1137 <autoTerminate>true</autoTerminate>
1138 <name>success</name>
1139 </relationships>
1140 <style>
1141 <entry>
1142 <key>background-color</key>
1143 <value>#6febb9</value>
1144 </entry>
1145 </style>
1146 <type>org.apache.nifi.processors.cassandra.PutCassandraQL</type>
1147 </processors>
1148 </snippet>
1149 <timestamp>08/29/2016 11:39:23 EDT</timestamp>
1150</template>