-
Notifications
You must be signed in to change notification settings - Fork 89
Expand file tree
/
Copy pathParquetOutputAdapterManager.h
More file actions
87 lines (62 loc) · 3.54 KB
/
Copy pathParquetOutputAdapterManager.h
File metadata and controls
87 lines (62 loc) · 3.54 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
#ifndef _IN_CSP_ADAPTERS_PARQUET_ParquetOutputAdapterManager_H
#define _IN_CSP_ADAPTERS_PARQUET_ParquetOutputAdapterManager_H
#include <csp/adapters/parquet/ParquetReader.h>
#include <csp/adapters/utils/StructAdapterInfo.h>
#include <csp/core/Generator.h>
#include <csp/core/Platform.h>
#include <csp/engine/AdapterManager.h>
#include <csp/engine/Dictionary.h>
#include <set>
#include <string>
#include <unordered_map>
#include <csp/adapters/parquet/DialectGenericListWriterInterface.h>
namespace csp::adapters::parquet
{
class ParquetWriter;
class ParquetOutputFilenameAdapter;
class ParquetDictBasketOutputWriter;
//Top level AdapterManager object for all parquet adapters in the engine
class CSPPARQUETADAPTERIMPL_EXPORT ParquetOutputAdapterManager final : public csp::AdapterManager
{
public:
using FileVisitorCallback = std::function<void(const std::string &)>;
ParquetOutputAdapterManager( csp::Engine *engine, const Dictionary &properties, FileVisitorCallback fileVisitor );
~ParquetOutputAdapterManager();
const char *name() const override{ return "ParquetOutputAdapterManager"; }
const std::string &getFileName() const{ return m_fileName; }
const std::string &getTimestampColumnName() const{ return m_timestampColumnName; }
bool isAllowOverwrite() const{ return m_allowOverwrite; }
uint32_t getBatchSize() const{ return m_batchSize; }
std::string getCompression() const{ return m_compression; }
bool isWriteArrowBinary() const{ return m_writeArrowBinary; }
bool isSplitColumnsToFiles() const{ return m_splitColumnsToFiles; }
//start the writer, open file if necessary
void start( DateTime starttime, DateTime endtime ) override;
//stop the writer, write any unwritten data and close file
void stop() override;
DateTime processNextSimTimeSlice( DateTime time ) override;
OutputAdapter *getOutputAdapter( CspTypePtr &type, const Dictionary &properties );
OutputAdapter *getListOutputAdapter( CspTypePtr &elemType, const Dictionary &properties,
const DialectGenericListWriterInterface::Ptr& listWriterInterface );
ParquetDictBasketOutputWriter *createDictOutputBasketWriter( const char *columnName, const CspTypePtr &cspTypePtr);
OutputAdapter *createOutputFileNameAdapter();
void changeFileName( const std::string &filename );
void scheduleEndCycle();
private:
OutputAdapter *getScalarOutputAdapter( CspTypePtr &type, const Dictionary &properties );
OutputAdapter *getStructOutputAdapter( CspTypePtr &type, const Dictionary &properties );
std::string m_fileName;
std::string m_timestampColumnName;
bool m_allowOverwrite;
uint32_t m_batchSize;
std::string m_compression;
bool m_writeArrowBinary;
bool m_splitColumnsToFiles;
std::unique_ptr<ParquetWriter> m_parquetWriter;
std::unordered_map<std::string, int> m_dictBasketWriterIndexByName;
std::vector<std::unique_ptr<ParquetDictBasketOutputWriter>> m_dictBasketWriters;
FileVisitorCallback m_fileVisitor;
ParquetOutputFilenameAdapter *m_outputFilenameAdapter;
};
}
#endif